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