Skip to main content

mj_controller/pollers/
runtime_feed.rs

1use super::*;
2use mj_client::runtime_feed::RuntimeProjection;
3use mj_core::snapshot_map::SnapshotMap;
4
5pub struct RemoteDashboardWorkerPoller {
6    pub updates: SessionManagerUpdates,
7    pub control: SessionManagerControl,
8    pub shutdown: SessionManagerShutdown,
9    pub state: tokio::sync::watch::Receiver<RuntimeStateUpdate>,
10    /// Reviews the daemon is running for this workspace's sessions.
11    pub reviews: tokio::sync::watch::Receiver<Vec<crate::review_host::RuntimeReviewView>>,
12    /// Background events the daemon wants reported once, oldest first.
13    pub notices: tokio::sync::watch::Receiver<Vec<daemon::RuntimeNotice>>,
14    /// What the daemon knows about quota. The daemon is the only prober.
15    pub quotas: tokio::sync::watch::Receiver<mj_client::quota::QuotaSnapshot>,
16    pub profile_capabilities:
17        tokio::sync::watch::Receiver<mj_core::profile_capabilities::ProfileCapabilitiesSnapshot>,
18    pub config: tokio::sync::watch::Receiver<mj_core::config::Config>,
19    pub health: tokio::sync::watch::Receiver<RuntimeFeedHealth>,
20}
21
22/// Records and lifecycle ownership must reach the surface in the same frame.
23#[derive(Debug, Clone, Default, PartialEq)]
24pub struct RuntimeStateUpdate {
25    pub session_cpu: SnapshotMap<String, mj_client::runtime_feed::SessionCpuView>,
26    pub last_subagent_policy: mj_core::subagent::SubagentPolicy,
27    pub native_agents: SnapshotMap<String, mj_core::native_agent::NativeAgentView>,
28    pub workspace_names: std::collections::BTreeMap<String, String>,
29    pub revision: u64,
30    pub records: SnapshotMap<String, SessionRecord>,
31    pub lifecycles: Vec<daemon::RuntimeLifecycleView>,
32    pub moves: SnapshotMap<String, mj_core::state::MoveOperation>,
33    /// Parent/child relations for the sessions in `records`, so a surface can
34    /// keep a daemon-created child out of the real workspace without a full
35    /// state reload.
36    pub subagents: SnapshotMap<String, SubagentRecord>,
37    /// What new-session defaults are chosen from; see
38    /// [`mj_client::runtime_feed::LaunchRecency`].
39    pub launch_recency: Vec<mj_client::runtime_feed::LaunchRecency>,
40    /// The daemon storage owner's verdict per target host.
41    pub storage: Vec<mj_core::targets::storage::TargetStorageView>,
42}
43
44/// What a session looked like the last time a view was published for it.
45///
46/// The poller compares this before reading anything, so a session that has not
47/// moved costs one comparison rather than a full transcript load. Nothing here
48/// grows with the transcript: the projection is identified by its ordinal and
49/// digest, and the operational state is bounded by the relay's own command and
50/// configuration surface.
51#[derive(Debug, Clone, PartialEq)]
52pub(super) struct PublishedView {
53    pub(super) projection_ordinal: u64,
54    pub(super) projection_digest: String,
55    pub(super) operational: Option<mj_core::relay::RelayOperationalState>,
56    pub(super) connected: bool,
57    pub(super) error: Option<String>,
58}
59
60impl PublishedView {
61    pub(super) fn of(runtime: &crate::daemon::RuntimeSessionView) -> Self {
62        Self {
63            projection_ordinal: runtime.projection_ordinal,
64            projection_digest: runtime.projection_digest.clone(),
65            operational: runtime.operational.clone(),
66            connected: runtime.connected,
67            error: runtime.error.as_ref().map(|error| format!("{error:?}")),
68        }
69    }
70
71    pub(super) fn matches(&self, runtime: &crate::daemon::RuntimeSessionView) -> bool {
72        *self == Self::of(runtime)
73    }
74}
75
76pub(super) const PROJECTION_CONVERGENCE_RETRIES: u8 = 20;
77pub(super) const PROJECTION_CONVERGENCE_RETRY_DELAY: Duration = Duration::from_millis(50);
78
79#[derive(Debug, Clone, PartialEq, Eq)]
80pub(super) struct ProjectionMismatch {
81    pub(super) published_ordinal: u64,
82    pub(super) published_digest: String,
83    pub(super) durable_ordinal: u64,
84    pub(super) durable_digest: String,
85}
86
87#[derive(Default)]
88pub(super) struct ProjectionConvergence {
89    pub(super) attempts: std::collections::BTreeMap<String, (ProjectionMismatch, u8)>,
90}
91
92impl ProjectionConvergence {
93    pub(super) fn converged(&mut self, session_id: &str) {
94        self.attempts.remove(session_id);
95    }
96
97    /// Give a lifecycle rollback and the daemon's cached relay view a bounded
98    /// window to converge. Repeating the same mismatch eventually reports the
99    /// integrity failure instead of hiding it indefinitely.
100    pub(super) fn should_retry(&mut self, session_id: &str, mismatch: ProjectionMismatch) -> bool {
101        let entry = self
102            .attempts
103            .entry(session_id.to_owned())
104            .or_insert_with(|| (mismatch.clone(), 0));
105        if entry.0 != mismatch {
106            *entry = (mismatch, 0);
107        }
108        entry.1 = entry.1.saturating_add(1);
109        entry.1 <= PROJECTION_CONVERGENCE_RETRIES
110    }
111}
112
113/// Read-only updates shared by the dashboard and workspace preview. A snapshot
114/// precedes its session views, so consumers can establish membership first.
115pub enum RuntimeFeedUpdate {
116    Snapshot(Box<RuntimeProjection>),
117    Session {
118        session_id: String,
119        view: Box<ManagedSessionView>,
120    },
121    /// The daemon no longer runs an actor for this session, so it has no view.
122    SessionRemoved(String),
123    Error(String),
124}
125
126/// Dropping a subscription cancels even a pending daemon long poll. The task
127/// owns no writer or relay connection; blocking projection reads are bounded.
128pub struct RuntimeFeed {
129    pub updates: tokio::sync::mpsc::Receiver<RuntimeFeedUpdate>,
130    pub(super) task: tokio::task::JoinHandle<()>,
131}
132
133impl Drop for RuntimeFeed {
134    fn drop(&mut self) {
135        self.task.abort();
136    }
137}
138
139pub(super) type StoredProjection = Option<(MaterializedSession, mj_core::state::ProjectionWindow)>;
140
141pub(super) static PROJECTION_READERS: std::sync::LazyLock<Arc<tokio::sync::Semaphore>> =
142    std::sync::LazyLock::new(|| Arc::new(tokio::sync::Semaphore::new(4)));
143
144/// A tail fetched at a cursor the daemon no longer holds. The feed takes a
145/// new snapshot and fetches again; nothing about the session is wrong.
146#[derive(Debug)]
147pub(super) struct TailCursorExpired;
148
149impl std::fmt::Display for TailCursorExpired {
150    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
151        f.write_str("the daemon no longer holds this runtime cursor")
152    }
153}
154
155impl std::error::Error for TailCursorExpired {}
156
157pub(super) fn spawn_runtime_feed_with<P, PF, L, LF, S>(
158    workspace_id: String,
159    poll: P,
160    load: L,
161) -> RuntimeFeed
162where
163    P: Fn(String, u64) -> PF + Send + 'static,
164    PF: Future<Output = Result<S>> + Send,
165    S: Into<RuntimeProjection> + Send + 'static,
166    L: Fn(String) -> LF + Clone + Send + 'static,
167    LF: Future<Output = Result<StoredProjection>> + Send + 'static,
168{
169    let (tx, updates) = tokio::sync::mpsc::channel(32);
170    let task = tokio::spawn(async move {
171        let result = run_runtime_feed(workspace_id, poll, load, &tx).await;
172        if let Err(error) = result {
173            let message = format!("Runtime feed stopped: {error:#}");
174            tracing::error!(%message);
175            let _ = tx.send(RuntimeFeedUpdate::Error(message)).await;
176        }
177    });
178    RuntimeFeed { updates, task }
179}
180
181/// A stopped daemon is a state the person can act on, not an I/O error.
182fn refresh_failure_notice(error: &anyhow::Error) -> String {
183    if daemon::daemon_not_running(error).is_some() {
184        "Mjolnir daemon is stopped; sessions refresh when it starts again.".to_owned()
185    } else {
186        format!("Could not refresh sessions: {error:#}")
187    }
188}
189
190pub(super) async fn run_runtime_feed<P, PF, L, LF, S>(
191    workspace_id: String,
192    poll: P,
193    load: L,
194    tx: &tokio::sync::mpsc::Sender<RuntimeFeedUpdate>,
195) -> Result<()>
196where
197    P: Fn(String, u64) -> PF,
198    PF: Future<Output = Result<S>>,
199    S: Into<RuntimeProjection>,
200    L: Fn(String) -> LF + Clone + Send + 'static,
201    LF: Future<Output = Result<StoredProjection>> + Send + 'static,
202{
203    let mut revision = 0;
204    let mut convergence = ProjectionConvergence::default();
205    let mut published = std::collections::BTreeMap::<String, PublishedView>::new();
206    let mut previous_sessions = SnapshotMap::new();
207    let mut previous_transcripts = SnapshotMap::new();
208    let mut pending_ids = std::collections::BTreeSet::new();
209    loop {
210        let mut snapshot = match poll(workspace_id.clone(), revision).await {
211            Ok(snapshot) => snapshot.into(),
212            Err(error) => {
213                if tx
214                    .send(RuntimeFeedUpdate::Error(refresh_failure_notice(&error)))
215                    .await
216                    .is_err()
217                {
218                    return Ok(());
219                }
220                tokio::time::sleep(Duration::from_millis(250)).await;
221                continue;
222            }
223        };
224        let snapshot_revision = snapshot.revision;
225        let sessions = std::mem::take(&mut snapshot.sessions);
226        let mut removed = Vec::new();
227        for (id, runtime) in previous_sessions.changes(&sessions) {
228            match runtime {
229                Some(runtime) if !published.get(id).is_some_and(|last| last.matches(runtime)) => {
230                    pending_ids.insert(id.clone());
231                }
232                None => {
233                    published.remove(id);
234                    pending_ids.remove(id);
235                    convergence.attempts.remove(id);
236                    removed.push(id.clone());
237                }
238                Some(_) => {}
239            }
240        }
241        previous_sessions = sessions.clone();
242        // A tail can change without its session's ordinal moving: a relay
243        // actor reloaded from the store reports compacted content this way.
244        for (id, tail) in previous_transcripts.changes(&snapshot.transcripts) {
245            if tail.is_some() && sessions.contains_key(id) {
246                published.remove(id);
247                pending_ids.insert(id.clone());
248            }
249        }
250        previous_transcripts = snapshot.transcripts.clone();
251        if tx
252            .send(RuntimeFeedUpdate::Snapshot(Box::new(snapshot)))
253            .await
254            .is_err()
255        {
256            return Ok(());
257        }
258        // Forgetting a view here is announced, so a consumer never keeps a
259        // copy this loop will not send again.
260        for session_id in removed {
261            if tx
262                .send(RuntimeFeedUpdate::SessionRemoved(session_id))
263                .await
264                .is_err()
265            {
266                return Ok(());
267            }
268        }
269        let mut pending = pending_ids
270            .iter()
271            .filter_map(|id| sessions.get(id).cloned())
272            .collect::<std::collections::VecDeque<_>>();
273        let mut tasks = tokio::task::JoinSet::new();
274        let mut retry = false;
275        while !pending.is_empty() || !tasks.is_empty() {
276            // Independent session reads overlap, without flooding SQLite or
277            // leaving an unbounded number of blocking reads after cancellation.
278            while tasks.len() < 4 {
279                let Some(runtime) = pending.pop_front() else {
280                    break;
281                };
282                let load = load.clone();
283                tasks.spawn(async move {
284                    let stored = if runtime.operational.is_some() {
285                        load(runtime.session_id.clone()).await
286                    } else {
287                        Ok(None)
288                    };
289                    (runtime, stored)
290                });
291            }
292            let Some(result) = tasks.join_next().await else {
293                break;
294            };
295            let (runtime, stored) = result.context("join runtime projection reader")?;
296            let session_id = runtime.session_id.clone();
297            let fingerprint = PublishedView::of(&runtime);
298            let expects_projection = runtime.operational.is_some();
299            let Some(view) = runtime_projection_view(runtime, stored, &mut convergence) else {
300                retry = true;
301                continue;
302            };
303            if view.snapshot.is_some() || !expects_projection {
304                published.insert(session_id.clone(), fingerprint);
305                pending_ids.remove(&session_id);
306            } else {
307                published.remove(&session_id);
308            }
309            if tx
310                .send(RuntimeFeedUpdate::Session {
311                    session_id,
312                    view: Box::new(view),
313                })
314                .await
315                .is_err()
316            {
317                return Ok(());
318            }
319        }
320        if retry {
321            tokio::time::sleep(PROJECTION_CONVERGENCE_RETRY_DELAY).await;
322        } else {
323            revision = revision.max(snapshot_revision);
324        }
325    }
326}
327
328pub(super) fn runtime_projection_view(
329    runtime: daemon::RuntimeSessionView,
330    stored: Result<StoredProjection>,
331    convergence: &mut ProjectionConvergence,
332) -> Option<ManagedSessionView> {
333    let Some(operational) = runtime.operational else {
334        return Some(ManagedSessionView {
335            snapshot: None,
336            connected: runtime.connected,
337            error: runtime.error,
338        });
339    };
340    let detail = match stored {
341        Ok(Some((materialized, window)))
342            if materialized.applied_event_ordinal > runtime.projection_ordinal
343                || (materialized.applied_event_ordinal == runtime.projection_ordinal
344                    && materialized.applied_event_digest == runtime.projection_digest) =>
345        {
346            convergence.converged(&runtime.session_id);
347            return Some(ManagedSessionView {
348                snapshot: Some(ManagedSessionSnapshot {
349                    materialized,
350                    window,
351                    operational,
352                    latest_credential_sync_signal: runtime.latest_credential_sync_signal,
353                    worker_build: None,
354                    subagent_requests: Vec::new(),
355                    subagent_results: Vec::new(),
356                }),
357                connected: runtime.connected,
358                error: runtime.error,
359            });
360        }
361        Ok(Some((materialized, _))) => {
362            let mismatch = ProjectionMismatch {
363                published_ordinal: runtime.projection_ordinal,
364                published_digest: runtime.projection_digest,
365                durable_ordinal: materialized.applied_event_ordinal,
366                durable_digest: materialized.applied_event_digest.clone(),
367            };
368            if convergence.should_retry(&runtime.session_id, mismatch) {
369                return None;
370            }
371            if materialized.applied_event_ordinal < runtime.projection_ordinal {
372                format!(
373                    "daemon published projection {} but SQLite contains only {} after a bounded convergence retry",
374                    runtime.projection_ordinal, materialized.applied_event_ordinal
375                )
376            } else {
377                format!(
378                    "daemon and SQLite projection digests differ at ordinal {} after a bounded convergence retry",
379                    runtime.projection_ordinal
380                )
381            }
382        }
383        Ok(None) => "daemon published a session with no transcript tail".into(),
384        Err(error) if error.is::<TailCursorExpired>() => return None,
385        Err(error) => format!("load daemon-owned projection: {error:#}"),
386    };
387    Some(ManagedSessionView {
388        snapshot: None,
389        connected: false,
390        error: Some(ViewError::ProjectionIntegrity(detail)),
391    })
392}
393
394#[cfg(test)]
395mod refresh_notice_tests {
396    use super::*;
397
398    #[test]
399    fn a_stopped_daemon_is_reported_as_stopped_not_as_a_file_error() {
400        let stopped = anyhow::Error::new(daemon::DaemonNotRunning {
401            metadata_path: "/data/daemon.json".into(),
402        });
403        assert_eq!(
404            refresh_failure_notice(&stopped),
405            "Mjolnir daemon is stopped; sessions refresh when it starts again."
406        );
407        assert_eq!(
408            refresh_failure_notice(&anyhow::anyhow!("unexpected end of file")),
409            "Could not refresh sessions: unexpected end of file"
410        );
411    }
412}