mj_controller/daemon/
subagent_park.rs1use super::*;
7
8const PARK_TIMEOUT: Duration = Duration::from_secs(120);
10
11const UNPARK_TIMEOUT: Duration = Duration::from_secs(600);
14
15impl RuntimeState {
16 const WAIT_PROMPT_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(1);
17
18 pub(crate) fn install_wait_prompt_backend(
19 &self,
20 backend: &Arc<crate::server_runtime::api::ApiBackend>,
21 ) {
22 assert!(
23 self.wait_prompt_backend
24 .set(Arc::downgrade(backend))
25 .is_ok(),
26 "delegation wait-prompt backend is installed once"
27 );
28 }
29
30 pub(crate) async fn ensure_parent_wait_prompt(&self, parent_session_id: &str) -> Result<()> {
31 if let Some(backend) = self
32 .wait_prompt_backend
33 .get()
34 .and_then(std::sync::Weak::upgrade)
35 {
36 match backend.ensure_parent_wait_prompt(parent_session_id).await {
37 Ok(()) => {
38 self.wait_prompt_retries
39 .lock()
40 .unwrap_or_else(PoisonError::into_inner)
41 .remove(parent_session_id);
42 }
43 Err(error) => {
44 self.schedule_wait_prompt_retry(parent_session_id);
45 return Err(error);
46 }
47 }
48 }
49 Ok(())
50 }
51
52 pub(crate) fn schedule_wait_prompt_retry(&self, parent_session_id: &str) {
53 self.wait_prompt_retries
54 .lock()
55 .unwrap_or_else(PoisonError::into_inner)
56 .insert(
57 parent_session_id.to_owned(),
58 std::time::Instant::now() + Self::WAIT_PROMPT_RETRY_DELAY,
59 );
60 }
61
62 pub(crate) fn take_due_wait_prompt_retries(
63 &self,
64 now: std::time::Instant,
65 limit: usize,
66 ) -> Vec<String> {
67 let mut retries = self
68 .wait_prompt_retries
69 .lock()
70 .unwrap_or_else(PoisonError::into_inner);
71 let due = retries
72 .iter()
73 .filter(|(_, due)| **due <= now)
74 .take(limit)
75 .map(|(parent, _)| parent.clone())
76 .collect::<Vec<_>>();
77 for parent in &due {
78 retries.remove(parent);
79 }
80 due
81 }
82
83 pub async fn park_subagent(
87 self: &Arc<Self>,
88 child_session_id: String,
89 ) -> Result<crate::controller::ParkOutcome> {
90 self.stop_idle_subagent(child_session_id, None).await
91 }
92
93 pub(super) async fn fail_subagent_start(self: &Arc<Self>, child_session_id: &str, cause: &str) {
100 let parent_session_id = self
101 .owner()
102 .controller()
103 .state
104 .subagents
105 .get(child_session_id)
106 .map(|relation| relation.parent_session_id.clone());
107 if !self
108 .owner()
109 .controller()
110 .state
111 .subagents
112 .contains_key(child_session_id)
113 {
114 return;
115 }
116 match self
117 .stop_idle_subagent(child_session_id.to_owned(), Some(cause.to_owned()))
118 .await
119 {
120 Ok(crate::controller::ParkOutcome::Parked) => {
121 tracing::info!(
122 session_id = child_session_id,
123 cause,
124 "sub-agent recorded as failed: its first prompt was refused"
125 );
126 }
127 Ok(outcome) => tracing::warn!(
128 session_id = child_session_id,
129 ?outcome,
130 "a sub-agent whose first prompt was refused was left running"
131 ),
132 Err(error) => tracing::warn!(
133 session_id = child_session_id,
134 error = format!("{error:#}"),
135 "could not record a sub-agent whose first prompt was refused as failed"
136 ),
137 }
138 if let Some(parent_session_id) = parent_session_id
139 && let Err(error) = self.ensure_parent_wait_prompt(&parent_session_id).await
140 {
141 tracing::warn!(
142 child_session_id,
143 parent_session_id,
144 error = format!("{error:#}"),
145 "could not reconcile the parent's sub-agent wait prompt after startup failure"
146 );
147 }
148 }
149
150 async fn stop_idle_subagent(
153 self: &Arc<Self>,
154 child_session_id: String,
155 failure: Option<String>,
156 ) -> Result<crate::controller::ParkOutcome> {
157 let result = self
158 .run_lifecycle(
159 child_session_id,
160 LifecycleKind::Park,
161 move |state, session_id, _cancelled| async move {
162 let never = AtomicBool::new(false);
166 let _recovery_reservation = blocking({
167 let observer = state.recovery_observer.clone();
168 let session_id = session_id.clone();
169 move || reserve_recovery_or_cancel(&observer, &session_id, &never)
170 })
171 .await?;
172 if state.close_is_requested(&session_id) {
173 return Ok(DaemonLifecycleResult::Park(
174 crate::controller::ParkOutcome::NotRunning,
175 ));
176 }
177 let controller = blocking(Controller::load).await?;
178 let executor = CancellableProcessExecutor::with_timeout(PARK_TIMEOUT);
179 let outcome = match &failure {
180 None => {
181 controller
182 .park_subagent_worker(
183 &session_id,
184 &executor,
185 &state.session_manager,
186 )
187 .await?
188 }
189 Some(cause) => {
190 controller
191 .fail_subagent_start_worker(
192 &session_id,
193 cause,
194 &executor,
195 &state.session_manager,
196 )
197 .await?
198 }
199 };
200 tracing::info!(%session_id, ?outcome, "sub-agent park finished");
201 Ok(DaemonLifecycleResult::Park(outcome))
202 },
203 )
204 .await?;
205 self.publish_revision();
206 match result {
207 DaemonLifecycleResult::Park(outcome) => Ok(outcome),
208 _ => unreachable!("a park returns its worker's reservation outcome"),
209 }
210 }
211
212 pub async fn unpark_subagent(self: &Arc<Self>, child_session_id: String) -> Result<()> {
216 self.wait_for_subagent_park(&child_session_id).await;
220 let parked = self
221 .session_state(&child_session_id)
222 .is_some_and(|state| state == SessionState::Parked);
223 if parked {
224 self.run_lifecycle(
225 child_session_id.clone(),
226 LifecycleKind::Unpark,
227 |state, session_id, cancelled| async move {
230 let _recovery_reservation = blocking({
231 let observer = state.recovery_observer.clone();
232 let session_id = session_id.clone();
233 let cancelled = cancelled.clone();
234 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
235 })
236 .await?;
237 let controller = blocking(Controller::load).await?;
238 let executor =
239 CancellableProcessExecutor::new(cancelled).with_deadline(UNPARK_TIMEOUT);
240 controller
241 .unpark_subagent_worker(&session_id, &executor)
242 .await?;
243 Ok(DaemonLifecycleResult::Done)
244 },
245 )
246 .await?;
247 self.publish_revision();
248 }
249 Ok(())
250 }
251
252 pub fn subagent_park_running(&self, child_session_id: &str) -> bool {
254 self.owner()
255 .lifecycle
256 .get(child_session_id)
257 .is_some_and(|active| active.kind == LifecycleKind::Park && active.is_running())
258 }
259
260 async fn wait_for_subagent_park(self: &Arc<Self>, child_session_id: &str) {
262 let pending = self
263 .owner()
264 .lifecycle
265 .get(child_session_id)
266 .filter(|active| active.kind == LifecycleKind::Park)
267 .map(|active| active.result.clone());
268 if let Some(pending) = pending {
269 if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
270 tracing::debug!(
271 session_id = %child_session_id,
272 error = format!("{error:#}"),
273 "the sub-agent park a restart waited for failed"
274 );
275 }
276 self.remove_completed_lifecycle(&pending);
277 }
278 }
279}