mj_controller/daemon/
snapshot.rs1use super::*;
2
3impl RuntimeState {
4 pub async fn reload_controller(&self) -> Result<()> {
5 let _mutation = self.config_mutation.lock().await;
8 let controller_loader = self.controller_loader;
9 let controller = tokio::task::spawn_blocking(controller_loader)
10 .await
11 .context("daemon controller reload task panicked")??;
12 let session_count = controller.state.sessions.len();
13 *self
14 .controller
15 .lock()
16 .unwrap_or_else(PoisonError::into_inner) = controller;
17 let revision = self.publish_revision();
18 tracing::debug!(revision, session_count, "daemon controller state reloaded");
19 Ok(())
20 }
21
22 pub(super) fn missing_target_record(
25 &self,
26 session_id: &str,
27 view: &ManagedSessionView,
28 ) -> Option<(String, String)> {
29 let Some(ViewError::TargetMissing(detail)) = &view.error else {
30 return None;
31 };
32 if view.connected {
33 return None;
34 }
35 if self
36 .lifecycle
37 .lock()
38 .unwrap_or_else(PoisonError::into_inner)
39 .get(session_id)
40 .is_some_and(|active| active.result.borrow().is_none())
41 {
42 return None;
43 }
44 let controller = self
45 .controller
46 .lock()
47 .unwrap_or_else(PoisonError::into_inner);
48 let session = controller.state.sessions.get(session_id)?;
49 if !matches!(
50 session.state,
51 SessionState::Running | SessionState::Disconnected
52 ) {
53 return None;
54 }
55 Some((detail.clone(), session.updated_at.clone()))
56 }
57
58 pub(super) async fn persist_missing_target(
59 &self,
60 session_id: &str,
61 detail: String,
62 observed_updated_at: String,
63 ) -> Result<()> {
64 let changed = blocking({
65 let session_id = session_id.to_owned();
66 let detail = detail.clone();
67 move || {
68 crate::database::mark_session_target_missing_if_current(
69 &session_id,
70 &detail,
71 &chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
72 &observed_updated_at,
73 )
74 }
75 })
76 .await?;
77 if changed.is_some() {
78 self.reload_controller().await?;
79 self.push_notice(session_id, detail);
80 }
81 Ok(())
82 }
83
84 pub(super) async fn publish_session(
85 &self,
86 session_id: String,
87 view: ManagedSessionView,
88 ) -> Result<()> {
89 let connected = view.connected;
90 let has_snapshot = view.snapshot.is_some();
91 tracing::debug!(
92 %session_id,
93 connected,
94 has_snapshot,
95 "daemon received a session view"
96 );
97 if let Some(snapshot) = view.snapshot.as_ref() {
98 let controller = self
99 .controller
100 .lock()
101 .unwrap_or_else(PoisonError::into_inner);
102 if let Some(session) = controller.state.sessions.get(&session_id).cloned() {
103 let quiet =
108 view.connected && snapshot.operational.safe_to_replace(session.harness_kind);
109 if quiet {
110 self.worker_upgrade_observer
111 .observe(WorkerUpgradeObservation {
112 session: session.clone(),
113 config: controller.config.clone(),
114 worker_build: snapshot.worker_build.clone(),
115 quiet,
116 });
117 }
118 self.recovery_observer.observe(RecoveryObservation {
119 checkpoint_safe: snapshot
120 .operational
121 .safe_for_checkpoint(session.harness_kind),
122 session,
123 config: controller.config.clone(),
124 latest_completed_turn_ordinal: snapshot.latest_completed_turn_ordinal(),
125 execution: snapshot.materialized.execution,
126 });
127 }
128 }
129 self.sessions
130 .lock()
131 .unwrap_or_else(PoisonError::into_inner)
132 .insert(
133 session_id.clone(),
134 RuntimeSessionView::from_managed(session_id, view),
135 );
136 reach_test_hook("relay_projection_before_revision_publication").await?;
137 self.publish_revision();
138 Ok(())
139 }
140
141 pub(super) async fn runtime_snapshot(
142 &self,
143 workspace_id: &str,
144 after_revision: u64,
145 all_workspaces: bool,
146 ) -> Result<RuntimeSnapshot> {
147 let mut revisions = self.revisions.subscribe();
148 if *revisions.borrow_and_update() <= after_revision {
149 let _ = tokio::time::timeout(Duration::from_secs(30), revisions.changed()).await;
150 }
151 let revision = self.revisions.current();
152 let moves = blocking(crate::database::load_move_operations).await?;
153 let workspace_names = blocking(crate::database::list_workspaces)
154 .await?
155 .into_iter()
156 .map(|workspace| (workspace.id, workspace.name))
157 .collect();
158 let session_ids = if all_workspaces {
159 self.controller
160 .lock()
161 .unwrap_or_else(PoisonError::into_inner)
162 .state
163 .sessions
164 .keys()
165 .cloned()
166 .collect::<BTreeSet<_>>()
167 } else {
168 let workspace_id = workspace_id.to_owned();
169 blocking(move || crate::database::session_ids_for_workspace(&workspace_id))
170 .await?
171 .into_iter()
172 .collect()
173 };
174 let sessions = self
175 .sessions
176 .lock()
177 .unwrap_or_else(PoisonError::into_inner)
178 .iter()
179 .filter(|(session_id, _)| session_ids.contains(*session_id))
180 .map(|(_, view)| view.clone())
181 .collect();
182 let controller = self
186 .controller
187 .lock()
188 .unwrap_or_else(PoisonError::into_inner);
189 let lifecycles = self
190 .lifecycle
191 .lock()
192 .unwrap_or_else(PoisonError::into_inner)
193 .iter()
194 .filter(|(session_id, active)| {
195 (all_workspaces
196 || session_ids.contains(*session_id)
197 || active.resume_workspace_id.as_deref() == Some(workspace_id))
198 && active.is_visible()
199 })
200 .map(|(session_id, active)| RuntimeLifecycleView {
201 operation_id: active.operation_id.clone(),
202 cancellable: active.is_cancellable()
203 && lifecycle_cancellable(
204 active.kind,
205 durable_session_state(&controller, session_id),
206 ),
207 session_id: session_id.clone(),
208 kind: active.kind.into(),
209 started_at_epoch_seconds: active.started_at_epoch_seconds,
210 active_stages: active
211 .active_stages
212 .iter()
213 .map(|(stage, (_, started_at))| (*stage, *started_at))
214 .collect(),
215 resume_destination: active.resume_destination.clone(),
216 notice: active.notice.clone(),
217 })
218 .collect();
219 let reviews = self
220 .review_host
221 .views()
222 .into_iter()
223 .filter(|review| session_ids.contains(&review.session_id))
224 .collect();
225 let notices = self
226 .notices
227 .lock()
228 .unwrap_or_else(PoisonError::into_inner)
229 .iter()
230 .filter(|notice| session_ids.contains(¬ice.session_id))
231 .cloned()
232 .collect();
233 let records = runtime_records_for_workspace(&controller, &session_ids);
234 let subagents = runtime_subagents_for_workspace(&controller, &records);
235 Ok(RuntimeSnapshot {
236 workspace_names,
237 moves: moves
238 .into_iter()
239 .filter(|operation| session_ids.contains(&operation.selection.session_id))
240 .collect(),
241 revision,
242 config: controller.config.clone(),
243 records,
244 sessions,
245 lifecycles,
246 reviews,
247 notices,
248 subagents,
249 })
250 }
251}