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 pub reviews: tokio::sync::watch::Receiver<Vec<crate::review_host::RuntimeReviewView>>,
12 pub notices: tokio::sync::watch::Receiver<Vec<daemon::RuntimeNotice>>,
14 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#[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 pub subagents: SnapshotMap<String, SubagentRecord>,
37 pub launch_recency: Vec<mj_client::runtime_feed::LaunchRecency>,
40}
41
42#[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 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
111pub enum RuntimeFeedUpdate {
114 Snapshot(Box<RuntimeProjection>),
115 Session {
116 session_id: String,
117 view: Box<ManagedSessionView>,
118 },
119 SessionRemoved(String),
121 Error(String),
122}
123
124pub 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#[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
179fn 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 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 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 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}