Skip to main content

mj_controller/pollers/
runtime_feed.rs

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