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
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 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}