1use std::time::Duration;
17
18use anyhow::{Context, Result, ensure};
19
20use mj_core::state::SessionState;
21
22use super::worker_binary::{
23 refresh_installed_worker_binary, replace_installed_worker_launch_config, stop_worker,
24 stop_worker_after_target_recovery,
25};
26use super::worker_restart::{InstalledWorkerRestart, WorkerRestartMessages};
27use super::{Controller, IdleWorkspaceLease};
28use crate::session_manager::SessionManagerControl;
29use crate::targets::{self, CommandExecutor};
30
31const PARK_ACTOR_TIMEOUT: Duration = Duration::from_secs(5);
33
34const PARK_RELEASE_TIMEOUT: Duration = Duration::from_secs(15);
37
38const RESTART_FROM_PARKED: WorkerRestartMessages = WorkerRestartMessages {
40 stop: "stop the parked sub-agent's worker",
41 replace: "install the current Mjolnir worker binary for the parked sub-agent",
42 start: "start the parked sub-agent's worker",
43 connect: "connect to the parked sub-agent's worker after starting it",
44 project_memory: "project memory will not be synchronized for the restarted sub-agent",
45 native_session: "wait for the parked sub-agent's harness to load its conversation",
46};
47
48#[derive(Debug, Clone, Copy, PartialEq, Eq)]
50pub enum ParkOutcome {
51 Parked,
53 Busy,
56 NotRunning,
59}
60
61fn cleanup_executor() -> targets::CancellableProcessExecutor {
64 targets::CancellableProcessExecutor::with_timeout(Duration::from_secs(30))
65}
66
67pub(super) fn explain_process_exhaustion(
75 error: anyhow::Error,
76 backend: &targets::TargetLocator,
77 session_id: &str,
78) -> anyhow::Error {
79 if !targets::shows_process_exhaustion(&error) {
80 return error;
81 }
82 let usage = targets::is_container(backend).then(|| {
83 let executor =
84 targets::CancellableProcessExecutor::with_timeout(targets::PIDS_USAGE_READ_TIMEOUT);
85 targets::read_pids_usage(&executor, backend, session_id)
86 });
87 anyhow::anyhow!(targets::process_exhaustion_message(usage.as_ref(), &error))
88}
89
90impl Controller {
91 pub async fn park_subagent_worker(
107 &self,
108 session_id: &str,
109 executor: &(impl CommandExecutor + Sync),
110 manager: &SessionManagerControl,
111 ) -> Result<ParkOutcome> {
112 ensure!(
113 self.state.subagents.contains_key(session_id),
114 "session {session_id} is not a sub-agent"
115 );
116 let Some(session) = self.state.sessions.get(session_id) else {
117 return Ok(ParkOutcome::NotRunning);
118 };
119 if session.state != SessionState::Running {
120 return Ok(ParkOutcome::NotRunning);
121 }
122 let (backend, worker_root) = self.worker_placement(session_id)?;
123 let handle = manager
124 .wait_for_session(session_id, PARK_ACTOR_TIMEOUT)
125 .await?;
126 let Some(mut lease) =
127 IdleWorkspaceLease::acquire_for_upgrade(&handle, session.harness_kind).await?
128 else {
129 return Ok(ParkOutcome::Busy);
130 };
131 if !lease.verify_for_upgrade().await? {
132 return Ok(ParkOutcome::Busy);
133 }
134 stop_worker_after_target_recovery(executor, &backend, session_id, &worker_root)
135 .context("stop the sub-agent's worker to park it")?;
136 let mut record = session.clone();
137 record.state = SessionState::Parked;
138 record.last_error = None;
139 record.updated_at = super::now();
140 crate::database::save_lifecycle_session(&record).context("record the parked sub-agent")?;
141 let released = tokio::time::timeout(PARK_RELEASE_TIMEOUT, async {
142 while manager.session(session_id.to_owned()).await.is_ok() {
143 tokio::time::sleep(Duration::from_millis(50)).await;
144 }
145 })
146 .await;
147 if released.is_err() {
148 tracing::warn!(
149 session_id,
150 "the session manager still held the parked sub-agent; releasing it anyway"
151 );
152 }
153 drop(lease);
154 Ok(ParkOutcome::Parked)
155 }
156
157 pub async fn unpark_subagent_worker(
172 &self,
173 session_id: &str,
174 executor: &(impl CommandExecutor + Sync),
175 ) -> Result<()> {
176 ensure!(
177 self.state.subagents.contains_key(session_id),
178 "session {session_id} is not a sub-agent"
179 );
180 let session = self
181 .state
182 .sessions
183 .get(session_id)
184 .with_context(|| format!("unknown session {session_id}"))?;
185 if session.state == SessionState::Running {
186 return Ok(());
187 }
188 ensure!(
189 session.state == SessionState::Parked,
190 "sub-agent {session_id} is {} and cannot be started again",
191 session.state.as_str()
192 );
193 let (backend, worker_root) = self.worker_placement(session_id)?;
194 let reconnect = targets::reconnect_plan(&backend, session_id)?
195 .commands
196 .into_iter()
197 .next()
198 .context("reconnect plan is empty")?;
199 let launch = self.current_worker_launch_config(session_id, &backend)?;
200 let started = async {
201 refresh_installed_worker_binary(executor, &backend, session_id)
202 .context(RESTART_FROM_PARKED.replace)?;
203 replace_installed_worker_launch_config(executor, &backend, session_id, &launch)
204 .context("install the current Mjolnir worker launch configuration")?;
205 let gate = super::provisioning::container_start_gate(&backend);
208 let _admitted = match &gate {
209 Some(gate) => gate.acquire().await.ok(),
210 None => None,
211 };
212 self.start_installed_worker(
213 session_id,
214 executor,
215 InstalledWorkerRestart {
216 backend: &backend,
217 worker_root: &worker_root,
218 reconnect: &reconnect,
219 launch: Some(&launch),
220 prepared: true,
221 messages: &RESTART_FROM_PARKED,
222 },
223 )
224 .await
225 }
226 .await;
227 let connection = match started {
228 Ok(connection) => connection,
229 Err(error) => {
230 if let Err(stop_error) = stop_worker(&cleanup_executor(), &backend, &worker_root) {
231 tracing::warn!(
232 session_id,
233 error = format!("{stop_error:#}"),
234 "could not stop the worker of a sub-agent whose restart failed"
235 );
236 }
237 return Err(explain_process_exhaustion(error, &backend, session_id));
238 }
239 };
240 drop(connection);
242 let mut record = session.clone();
243 record.state = SessionState::Running;
244 record.last_error = None;
245 record.updated_at = super::now();
246 if let Err(error) = crate::database::save_lifecycle_session(&record) {
247 if let Err(stop_error) = stop_worker(&cleanup_executor(), &backend, &worker_root) {
248 tracing::warn!(
249 session_id,
250 error = format!("{stop_error:#}"),
251 "could not stop the worker of a sub-agent whose restart was not recorded"
252 );
253 }
254 return Err(error.context("record the restarted sub-agent as running"));
255 }
256 Ok(())
257 }
258}
259
260#[cfg(all(test, unix))]
261mod tests {
262 use std::sync::Mutex;
263
264 use agent_client_protocol::schema::v1::ContentBlock;
265 use mj_core::relay::RelayCommand;
266 use mj_core::state::TargetLocator;
267
268 use super::*;
269 use crate::controller::checkpoint::tests::{
270 LATCH_RELAY_SESSION, ReleaseSupport, latch_relay_target,
271 };
272 use crate::controller::test_support::{IsolatedTest, checkpoint_test_session, test_name};
273 use crate::targets::{CommandOutput, CommandSpec};
274
275 const MARKER: &str = "MJ_TEST_SUBAGENT_PARK_CHILD";
276
277 fn isolated(test: &str) -> bool {
280 if std::env::var_os(MARKER).is_some() {
281 return true;
282 }
283 let directory = tempfile::tempdir().unwrap();
284 IsolatedTest::new(test_name(module_path!(), test))
285 .env(MARKER, "1")
286 .isolated_store(directory.path())
287 .run();
288 false
289 }
290
291 #[derive(Default)]
295 struct RacingStop {
296 purposes: Mutex<Vec<String>>,
297 racer: Mutex<
298 Option<(
299 crate::session_manager::ManagedSessionHandle,
300 tokio::runtime::Handle,
301 )>,
302 >,
303 raced: Mutex<Option<tokio::task::JoinHandle<Result<u64>>>>,
304 }
305
306 impl CommandExecutor for RacingStop {
307 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
308 self.purposes.lock().unwrap().push(command.purpose.clone());
309 if let Some((handle, runtime)) = self.racer.lock().unwrap().take() {
310 *self.raced.lock().unwrap() = Some(runtime.spawn(async move {
311 handle
312 .submit(
313 "raced-prompt".into(),
314 RelayCommand::Prompt {
315 prompt: vec![ContentBlock::from("one more thing")],
316 },
317 )
318 .await
319 }));
320 }
321 Ok(CommandOutput {
322 status: 0,
323 stdout: Vec::new(),
324 stderr: Vec::new(),
325 })
326 }
327 }
328
329 fn register_child(root: &std::path::Path) {
332 crate::database::save_session(&checkpoint_test_session("parent-1")).unwrap();
333 let mut child = checkpoint_test_session(LATCH_RELAY_SESSION);
334 child.target = Some(TargetLocator::LocalBare {
335 worker_root: root.join(LATCH_RELAY_SESSION),
336 });
337 crate::database::save_subagent_session(
338 &child,
339 &mj_core::subagent::SubagentRecord {
340 child_session_id: LATCH_RELAY_SESSION.into(),
341 parent_session_id: "parent-1".into(),
342 task_name: "map the parser".into(),
343 profile_id: "codex".into(),
344 model: None,
345 effort: None,
346 working_directory: Default::default(),
347 initial_prompt: "map the parser".into(),
348 request_key: "request-1".into(),
349 created_at: "2026-09-25T00:00:00Z".into(),
350 noticed_turn: None,
351 handback_tool: true,
352 },
353 )
354 .unwrap();
355 assert!(
356 crate::database::record_subagent_handback(
357 LATCH_RELAY_SESSION,
358 &mj_core::subagent::SubagentHandback {
359 command_id: "task-1".into(),
360 message: "The parser has three entry points.".into(),
361 recorded_at_ms: 1,
362 },
363 )
364 .unwrap()
365 );
366 }
367
368 fn loaded_controller() -> Controller {
369 Controller {
370 config: mj_core::config::Config::default(),
371 state: crate::database::load_state().unwrap(),
372 }
373 }
374
375 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
376 async fn parking_stops_an_idle_child_keeps_its_record_and_turns_away_a_racing_prompt() {
377 if !isolated("parking_stops_an_idle_child_keeps_its_record_and_turns_away_a_racing_prompt")
378 {
379 return;
380 }
381 let _writer = crate::database::install_isolated_test_writer();
382 let root = tempfile::tempdir().unwrap();
383 register_child(root.path());
384 let crate::session_manager::SessionManagerChannels {
385 targets,
386 control,
387 updates: _updates,
388 shutdown,
389 } = crate::session_manager::spawn_session_manager().unwrap();
390 targets
391 .send(vec![latch_relay_target(
392 root.path(),
393 None,
394 ReleaseSupport::Supported,
395 false,
396 )])
397 .unwrap();
398 let handle = control
399 .wait_for_session(LATCH_RELAY_SESSION, Duration::from_secs(10))
400 .await
401 .unwrap();
402 let executor = RacingStop::default();
403 *executor.racer.lock().unwrap() = Some((handle, tokio::runtime::Handle::current()));
404 let refresher = tokio::spawn(async move {
407 while crate::database::load_session_state(LATCH_RELAY_SESSION).unwrap()
408 != Some(SessionState::Parked)
409 {
410 tokio::time::sleep(Duration::from_millis(25)).await;
411 }
412 targets.send_replace(Vec::new());
413 targets
414 });
415
416 let outcome = loaded_controller()
417 .park_subagent_worker(LATCH_RELAY_SESSION, &executor, &control)
418 .await
419 .unwrap();
420
421 assert_eq!(outcome, ParkOutcome::Parked);
422 assert!(
423 executor
424 .purposes
425 .lock()
426 .unwrap()
427 .iter()
428 .any(|purpose| purpose == "stop Mjolnir worker daemon"),
429 "the child's worker was stopped: {:?}",
430 executor.purposes.lock().unwrap()
431 );
432 let raced = executor.raced.lock().unwrap().take();
436 let raced = raced
437 .expect("the stop raced a prompt")
438 .await
439 .unwrap()
440 .expect_err("a prompt that arrived during the park is turned away");
441 assert!(
442 raced
443 .downcast_ref::<mj_client::session::DeliveryUnconfirmed>()
444 .is_none(),
445 "a turned-away prompt is known not to be delivered: {raced:#}"
446 );
447 let stored = crate::database::load_state().unwrap();
449 let child = &stored.sessions[LATCH_RELAY_SESSION];
450 assert_eq!(child.state, SessionState::Parked);
451 assert!(child.target.is_some(), "a parked child keeps its target");
452 assert!(stored.subagents.contains_key(LATCH_RELAY_SESSION));
453 assert_eq!(
454 crate::database::load_subagent_report(LATCH_RELAY_SESSION)
455 .unwrap()
456 .handback
457 .map(|handback| handback.message)
458 .as_deref(),
459 Some("The parser has three entry points.")
460 );
461 assert!(control.session(LATCH_RELAY_SESSION).await.is_err());
462 let _targets = refresher.await.unwrap();
463 shutdown.shutdown().await.unwrap();
464 }
465
466 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
467 async fn a_child_with_work_in_flight_is_not_parked() {
468 if !isolated("a_child_with_work_in_flight_is_not_parked") {
469 return;
470 }
471 let _writer = crate::database::install_isolated_test_writer();
472 let root = tempfile::tempdir().unwrap();
473 register_child(root.path());
474 let channels = crate::session_manager::spawn_session_manager().unwrap();
475 channels
476 .targets
477 .send(vec![latch_relay_target(
478 root.path(),
479 None,
480 ReleaseSupport::Supported,
481 true,
482 )])
483 .unwrap();
484 channels
485 .control
486 .wait_for_session(LATCH_RELAY_SESSION, Duration::from_secs(10))
487 .await
488 .unwrap();
489 let executor = RacingStop::default();
490
491 let outcome = loaded_controller()
492 .park_subagent_worker(LATCH_RELAY_SESSION, &executor, &channels.control)
493 .await
494 .unwrap();
495
496 assert_eq!(outcome, ParkOutcome::Busy);
497 assert!(
498 executor.purposes.lock().unwrap().is_empty(),
499 "nothing stopped"
500 );
501 assert_eq!(
502 crate::database::load_session_state(LATCH_RELAY_SESSION).unwrap(),
503 Some(SessionState::Running)
504 );
505 channels.shutdown.shutdown().await.unwrap();
506 }
507
508 #[test]
509 fn only_a_full_target_is_rewritten_and_a_bare_one_reads_no_container_counts() {
510 let backend = targets::TargetLocator::LocalBare {
511 worker_root: "/tmp/workers/child".into(),
512 };
513 let unrelated =
514 explain_process_exhaustion(anyhow::anyhow!("the harness exited"), &backend, "child");
515 assert_eq!(format!("{unrelated:#}"), "the harness exited");
516
517 let full = explain_process_exhaustion(
518 anyhow::anyhow!("sh: 1: Cannot fork").context("start the parked sub-agent's worker"),
519 &backend,
520 "child",
521 );
522 let message = format!("{full:#}");
523 assert!(
524 message.starts_with("the target machine ran out of process slots"),
525 "{message}"
526 );
527 assert!(
528 message.contains("Close sub-agents you no longer need"),
529 "{message}"
530 );
531 assert!(
532 message.contains("Cannot fork"),
533 "the original error stays: {message}"
534 );
535 }
536}