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