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