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 config: tokio::sync::watch::Receiver<mj_core::config::Config>,
17 pub health: tokio::sync::watch::Receiver<RuntimeFeedHealth>,
18}
19
20#[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 pub subagents: SnapshotMap<String, SubagentRecord>,
34 pub launch_recency: Vec<mj_client::runtime_feed::LaunchRecency>,
37}
38
39#[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 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
108pub enum RuntimeFeedUpdate {
111 Snapshot(Box<RuntimeProjection>),
112 Session {
113 session_id: String,
114 view: Box<ManagedSessionView>,
115 },
116 SessionRemoved(String),
118 Error(String),
119}
120
121pub 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 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
186fn 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 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 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}