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 pub reviews: tokio::sync::watch::Receiver<Vec<crate::review_host::RuntimeReviewView>>,
11 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#[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 pub subagents: Vec<SubagentRecord>,
30}
31
32#[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 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
101pub enum RuntimeFeedUpdate {
104 Snapshot(Box<daemon::RuntimeSnapshot>),
105 Session {
106 session_id: String,
107 view: Box<ManagedSessionView>,
108 },
109 Error(String),
110}
111
112pub 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 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
176fn 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 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}