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
176pub(super) async fn run_runtime_feed<P, PF, L, LF>(
177    workspace_id: String,
178    poll: P,
179    load: L,
180    tx: &tokio::sync::mpsc::Sender<RuntimeFeedUpdate>,
181) -> Result<()>
182where
183    P: Fn(String, u64) -> PF,
184    PF: Future<Output = Result<daemon::RuntimeSnapshot>>,
185    L: Fn(String) -> LF + Clone + Send + 'static,
186    LF: Future<Output = Result<StoredProjection>> + Send + 'static,
187{
188    let mut revision = 0;
189    let mut convergence = ProjectionConvergence::default();
190    let mut published = std::collections::BTreeMap::<String, PublishedView>::new();
191    loop {
192        let mut snapshot = match poll(workspace_id.clone(), revision).await {
193            Ok(snapshot) => snapshot,
194            Err(error) => {
195                if tx
196                    .send(RuntimeFeedUpdate::Error(format!(
197                        "Could not refresh sessions: {error:#}"
198                    )))
199                    .await
200                    .is_err()
201                {
202                    return Ok(());
203                }
204                tokio::time::sleep(Duration::from_millis(250)).await;
205                continue;
206            }
207        };
208        let snapshot_revision = snapshot.revision;
209        let sessions = std::mem::take(&mut snapshot.sessions);
210        published.retain(|id, _| sessions.iter().any(|session| &session.session_id == id));
211        convergence
212            .attempts
213            .retain(|id, _| sessions.iter().any(|session| &session.session_id == id));
214        if tx
215            .send(RuntimeFeedUpdate::Snapshot(Box::new(snapshot)))
216            .await
217            .is_err()
218        {
219            return Ok(());
220        }
221        let mut pending = sessions
222            .into_iter()
223            .filter(|runtime| {
224                !published
225                    .get(&runtime.session_id)
226                    .is_some_and(|last| last.matches(runtime))
227            })
228            .collect::<std::collections::VecDeque<_>>();
229        let mut tasks = tokio::task::JoinSet::new();
230        let mut retry = false;
231        while !pending.is_empty() || !tasks.is_empty() {
232            // Independent session reads overlap, without flooding SQLite or
233            // leaving an unbounded number of blocking reads after cancellation.
234            while tasks.len() < 4 {
235                let Some(runtime) = pending.pop_front() else {
236                    break;
237                };
238                let load = load.clone();
239                tasks.spawn(async move {
240                    let stored = if runtime.operational.is_some() {
241                        load(runtime.session_id.clone()).await
242                    } else {
243                        Ok(None)
244                    };
245                    (runtime, stored)
246                });
247            }
248            let Some(result) = tasks.join_next().await else {
249                break;
250            };
251            let (runtime, stored) = result.context("join runtime projection reader")?;
252            let session_id = runtime.session_id.clone();
253            let fingerprint = PublishedView::of(&runtime);
254            let Some(view) = runtime_projection_view(runtime, stored, &mut convergence) else {
255                retry = true;
256                continue;
257            };
258            if view.snapshot.is_some() {
259                published.insert(session_id.clone(), fingerprint);
260            } else {
261                published.remove(&session_id);
262            }
263            if tx
264                .send(RuntimeFeedUpdate::Session {
265                    session_id,
266                    view: Box::new(view),
267                })
268                .await
269                .is_err()
270            {
271                return Ok(());
272            }
273        }
274        if retry {
275            tokio::time::sleep(PROJECTION_CONVERGENCE_RETRY_DELAY).await;
276        } else {
277            revision = revision.max(snapshot_revision);
278        }
279    }
280}
281
282pub(super) fn runtime_projection_view(
283    runtime: daemon::RuntimeSessionView,
284    stored: Result<StoredProjection>,
285    convergence: &mut ProjectionConvergence,
286) -> Option<ManagedSessionView> {
287    let Some(operational) = runtime.operational else {
288        return Some(ManagedSessionView {
289            snapshot: None,
290            connected: runtime.connected,
291            error: runtime.error,
292        });
293    };
294    let detail = match stored {
295        Ok(Some((materialized, window)))
296            if materialized.applied_event_ordinal > runtime.projection_ordinal
297                || (materialized.applied_event_ordinal == runtime.projection_ordinal
298                    && materialized.applied_event_digest == runtime.projection_digest) =>
299        {
300            convergence.converged(&runtime.session_id);
301            return Some(ManagedSessionView {
302                snapshot: Some(ManagedSessionSnapshot {
303                    materialized,
304                    window,
305                    operational,
306                    latest_credential_sync_signal: runtime.latest_credential_sync_signal,
307                    worker_build: None,
308                    subagent_requests: Vec::new(),
309                    subagent_results: Vec::new(),
310                }),
311                connected: runtime.connected,
312                error: runtime.error,
313            });
314        }
315        Ok(Some((materialized, _))) => {
316            let mismatch = ProjectionMismatch {
317                published_ordinal: runtime.projection_ordinal,
318                published_digest: runtime.projection_digest,
319                durable_ordinal: materialized.applied_event_ordinal,
320                durable_digest: materialized.applied_event_digest.clone(),
321            };
322            if convergence.should_retry(&runtime.session_id, mismatch) {
323                return None;
324            }
325            if materialized.applied_event_ordinal < runtime.projection_ordinal {
326                format!(
327                    "daemon published projection {} but SQLite contains only {} after a bounded convergence retry",
328                    runtime.projection_ordinal, materialized.applied_event_ordinal
329                )
330            } else {
331                format!(
332                    "daemon and SQLite projection digests differ at ordinal {} after a bounded convergence retry",
333                    runtime.projection_ordinal
334                )
335            }
336        }
337        Ok(None) => "daemon published a session with no durable projection".into(),
338        Err(error) => format!("load daemon-owned projection: {error:#}"),
339    };
340    Some(ManagedSessionView {
341        snapshot: None,
342        connected: false,
343        error: Some(ViewError::ProjectionIntegrity(detail)),
344    })
345}