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 self.stop_idle_subagent_worker(session_id, executor, manager, None)
113 .await
114 }
115
116 pub async fn fail_subagent_start_worker(
127 &self,
128 session_id: &str,
129 cause: &str,
130 executor: &(impl CommandExecutor + Sync),
131 manager: &SessionManagerControl,
132 ) -> Result<ParkOutcome> {
133 self.stop_idle_subagent_worker(session_id, executor, manager, Some(cause))
134 .await
135 }
136
137 async fn stop_idle_subagent_worker(
140 &self,
141 session_id: &str,
142 executor: &(impl CommandExecutor + Sync),
143 manager: &SessionManagerControl,
144 failure: Option<&str>,
145 ) -> Result<ParkOutcome> {
146 ensure!(
147 self.state.subagents.contains_key(session_id),
148 "session {session_id} is not a sub-agent"
149 );
150 let Some(session) = self.state.sessions.get(session_id) else {
151 return Ok(ParkOutcome::NotRunning);
152 };
153 if session.state != SessionState::Running {
154 return Ok(ParkOutcome::NotRunning);
155 }
156 let (backend, worker_root) = self.worker_placement(session_id)?;
157 let handle = manager
158 .wait_for_session(session_id, PARK_ACTOR_TIMEOUT)
159 .await?;
160 let Some(mut lease) =
161 IdleWorkspaceLease::acquire_for_upgrade(&handle, session.harness_kind).await?
162 else {
163 return Ok(ParkOutcome::Busy);
164 };
165 if !lease.verify_for_upgrade().await? {
166 return Ok(ParkOutcome::Busy);
167 }
168 stop_worker_after_target_recovery(executor, &backend, session_id, &worker_root)
169 .context("stop the sub-agent's worker")?;
170 let mut record = session.clone();
171 record.state = if failure.is_some() {
172 SessionState::Error
173 } else {
174 SessionState::Parked
175 };
176 record.last_error = failure.map(str::to_owned);
177 record.updated_at = super::now();
178 crate::database::save_lifecycle_session(&record).context("record the stopped sub-agent")?;
179 let released = tokio::time::timeout(PARK_RELEASE_TIMEOUT, async {
180 while manager.session(session_id.to_owned()).await.is_ok() {
181 tokio::time::sleep(Duration::from_millis(50)).await;
182 }
183 })
184 .await;
185 if released.is_err() {
186 tracing::warn!(
187 session_id,
188 "the session manager still held the stopped sub-agent; releasing it anyway"
189 );
190 }
191 drop(lease);
192 Ok(ParkOutcome::Parked)
193 }
194
195 pub async fn unpark_subagent_worker(
210 &self,
211 session_id: &str,
212 executor: &(impl CommandExecutor + Sync),
213 ) -> Result<()> {
214 ensure!(
215 self.state.subagents.contains_key(session_id),
216 "session {session_id} is not a sub-agent"
217 );
218 let session = self
219 .state
220 .sessions
221 .get(session_id)
222 .with_context(|| format!("unknown session {session_id}"))?;
223 if session.state == SessionState::Running {
224 return Ok(());
225 }
226 ensure!(
227 session.state == SessionState::Parked,
228 "sub-agent {session_id} is {} and cannot be started again",
229 session.state.as_str()
230 );
231 let (backend, worker_root) = self.worker_placement(session_id)?;
232 let reconnect = targets::reconnect_plan(&backend, session_id)?
233 .commands
234 .into_iter()
235 .next()
236 .context("reconnect plan is empty")?;
237 let launch = self.current_worker_launch_config(session_id, &backend)?;
238 let started = async {
239 refresh_installed_worker_binary(executor, &backend, session_id)
240 .context(RESTART_FROM_PARKED.replace)?;
241 replace_installed_worker_launch_config(executor, &backend, session_id, &launch)
242 .context("install the current Mjolnir worker launch configuration")?;
243 let gate = super::provisioning::container_start_gate(&backend);
246 let _admitted = match &gate {
247 Some(gate) => gate.acquire().await.ok(),
248 None => None,
249 };
250 self.start_installed_worker(
251 session_id,
252 executor,
253 InstalledWorkerRestart {
254 backend: &backend,
255 worker_root: &worker_root,
256 reconnect: &reconnect,
257 launch: Some(&launch),
258 prepared: true,
259 messages: &RESTART_FROM_PARKED,
260 },
261 )
262 .await
263 }
264 .await;
265 let connection = match started {
266 Ok(connection) => connection,
267 Err(error) => {
268 if let Err(stop_error) = stop_worker(&cleanup_executor(), &backend, &worker_root) {
269 tracing::warn!(
270 session_id,
271 error = format!("{stop_error:#}"),
272 "could not stop the worker of a sub-agent whose restart failed"
273 );
274 }
275 return Err(explain_process_exhaustion(error, &backend, session_id));
276 }
277 };
278 drop(connection);
280 let mut record = session.clone();
281 record.state = SessionState::Running;
282 record.last_error = None;
283 record.updated_at = super::now();
284 if let Err(error) = crate::database::save_lifecycle_session(&record) {
285 if let Err(stop_error) = stop_worker(&cleanup_executor(), &backend, &worker_root) {
286 tracing::warn!(
287 session_id,
288 error = format!("{stop_error:#}"),
289 "could not stop the worker of a sub-agent whose restart was not recorded"
290 );
291 }
292 return Err(error.context("record the restarted sub-agent as running"));
293 }
294 Ok(())
295 }
296}
297
298#[cfg(all(test, unix))]
299mod tests {
300 use std::sync::Mutex;
301
302 use agent_client_protocol::schema::v1::ContentBlock;
303 use mj_core::relay::RelayCommand;
304 use mj_core::state::TargetLocator;
305
306 use super::*;
307 use crate::controller::checkpoint::tests::{
308 LATCH_RELAY_SESSION, ReleaseSupport, latch_relay_target,
309 };
310 use crate::controller::test_support::{IsolatedTest, checkpoint_test_session, test_name};
311 use crate::targets::{CommandOutput, CommandSpec};
312
313 const MARKER: &str = "MJ_TEST_SUBAGENT_PARK_CHILD";
314
315 fn isolated(test: &str) -> bool {
318 if std::env::var_os(MARKER).is_some() {
319 return true;
320 }
321 let directory = tempfile::tempdir().unwrap();
322 IsolatedTest::new(test_name(module_path!(), test))
323 .env(MARKER, "1")
324 .isolated_store(directory.path())
325 .run();
326 false
327 }
328
329 #[derive(Default)]
333 struct RacingStop {
334 purposes: Mutex<Vec<String>>,
335 racer: Mutex<
336 Option<(
337 crate::session_manager::ManagedSessionHandle,
338 tokio::runtime::Handle,
339 )>,
340 >,
341 raced: Mutex<Option<tokio::task::JoinHandle<Result<u64>>>>,
342 }
343
344 impl CommandExecutor for RacingStop {
345 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
346 self.purposes.lock().unwrap().push(command.purpose.clone());
347 if let Some((handle, runtime)) = self.racer.lock().unwrap().take() {
348 *self.raced.lock().unwrap() = Some(runtime.spawn(async move {
349 handle
350 .submit(
351 "raced-prompt".into(),
352 RelayCommand::Prompt {
353 prompt: vec![ContentBlock::from("one more thing")],
354 },
355 )
356 .await
357 }));
358 }
359 Ok(CommandOutput {
360 status: 0,
361 stdout: Vec::new(),
362 stderr: Vec::new(),
363 })
364 }
365 }
366
367 fn register_child(root: &std::path::Path) {
370 crate::database::save_session(&checkpoint_test_session("parent-1")).unwrap();
371 let mut child = checkpoint_test_session(LATCH_RELAY_SESSION);
372 child.target = Some(TargetLocator::LocalBare {
373 worker_root: root.join(LATCH_RELAY_SESSION),
374 });
375 crate::database::save_subagent_session(
376 &child,
377 &mj_core::subagent::SubagentRecord {
378 child_session_id: LATCH_RELAY_SESSION.into(),
379 parent_session_id: "parent-1".into(),
380 task_name: "map the parser".into(),
381 profile_id: "codex".into(),
382 model: None,
383 effort: None,
384 working_directory: Default::default(),
385 initial_prompt: "map the parser".into(),
386 request_key: "request-1".into(),
387 created_at: "2026-09-25T00:00:00Z".into(),
388 noticed_turn: None,
389 handback_tool: true,
390 },
391 )
392 .unwrap();
393 assert!(
394 crate::database::record_subagent_handback(
395 LATCH_RELAY_SESSION,
396 &mj_core::subagent::SubagentHandback {
397 command_id: "task-1".into(),
398 message: "The parser has three entry points.".into(),
399 recorded_at_ms: 1,
400 },
401 )
402 .unwrap()
403 );
404 }
405
406 fn loaded_controller() -> Controller {
407 Controller {
408 config: mj_core::config::Config::default(),
409 state: crate::database::load_state().unwrap(),
410 }
411 }
412
413 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
414 async fn parking_stops_an_idle_child_keeps_its_record_and_turns_away_a_racing_prompt() {
415 if !isolated("parking_stops_an_idle_child_keeps_its_record_and_turns_away_a_racing_prompt")
416 {
417 return;
418 }
419 let _writer = crate::database::install_isolated_test_writer();
420 let root = tempfile::tempdir().unwrap();
421 register_child(root.path());
422 let crate::session_manager::SessionManagerChannels {
423 targets,
424 control,
425 updates: _updates,
426 shutdown,
427 } = crate::session_manager::spawn_session_manager().unwrap();
428 targets
429 .send(vec![latch_relay_target(
430 root.path(),
431 None,
432 ReleaseSupport::Supported,
433 false,
434 )])
435 .unwrap();
436 let handle = control
437 .wait_for_session(LATCH_RELAY_SESSION, Duration::from_secs(10))
438 .await
439 .unwrap();
440 let executor = RacingStop::default();
441 *executor.racer.lock().unwrap() = Some((handle, tokio::runtime::Handle::current()));
442 let refresher = tokio::spawn(async move {
445 while crate::database::load_session_state(LATCH_RELAY_SESSION).unwrap()
446 != Some(SessionState::Parked)
447 {
448 tokio::time::sleep(Duration::from_millis(25)).await;
449 }
450 targets.send_replace(Vec::new());
451 targets
452 });
453
454 let outcome = loaded_controller()
455 .park_subagent_worker(LATCH_RELAY_SESSION, &executor, &control)
456 .await
457 .unwrap();
458
459 assert_eq!(outcome, ParkOutcome::Parked);
460 assert!(
461 executor
462 .purposes
463 .lock()
464 .unwrap()
465 .iter()
466 .any(|purpose| purpose == "stop Mjolnir worker daemon"),
467 "the child's worker was stopped: {:?}",
468 executor.purposes.lock().unwrap()
469 );
470 let raced = executor.raced.lock().unwrap().take();
474 let raced = raced
475 .expect("the stop raced a prompt")
476 .await
477 .unwrap()
478 .expect_err("a prompt that arrived during the park is turned away");
479 assert!(
480 raced
481 .downcast_ref::<mj_client::session::DeliveryUnconfirmed>()
482 .is_none(),
483 "a turned-away prompt is known not to be delivered: {raced:#}"
484 );
485 let stored = crate::database::load_state().unwrap();
487 let child = &stored.sessions[LATCH_RELAY_SESSION];
488 assert_eq!(child.state, SessionState::Parked);
489 assert!(child.target.is_some(), "a parked child keeps its target");
490 assert!(stored.subagents.contains_key(LATCH_RELAY_SESSION));
491 assert_eq!(
492 crate::database::load_subagent_report(LATCH_RELAY_SESSION)
493 .unwrap()
494 .handback
495 .map(|handback| handback.message)
496 .as_deref(),
497 Some("The parser has three entry points.")
498 );
499 assert!(control.session(LATCH_RELAY_SESSION).await.is_err());
500 let _targets = refresher.await.unwrap();
501 shutdown.shutdown().await.unwrap();
502 }
503
504 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
508 async fn a_child_whose_start_failed_is_stopped_and_recorded_as_failed() {
509 if !isolated("a_child_whose_start_failed_is_stopped_and_recorded_as_failed") {
510 return;
511 }
512 let _writer = crate::database::install_isolated_test_writer();
513 let root = tempfile::tempdir().unwrap();
514 register_child(root.path());
515 let crate::session_manager::SessionManagerChannels {
516 targets,
517 control,
518 updates: _updates,
519 shutdown,
520 } = crate::session_manager::spawn_session_manager().unwrap();
521 targets
522 .send(vec![latch_relay_target(
523 root.path(),
524 None,
525 ReleaseSupport::Supported,
526 false,
527 )])
528 .unwrap();
529 control
530 .wait_for_session(LATCH_RELAY_SESSION, Duration::from_secs(10))
531 .await
532 .unwrap();
533 let refresher = tokio::spawn(async move {
534 while crate::database::load_session_state(LATCH_RELAY_SESSION).unwrap()
535 != Some(SessionState::Error)
536 {
537 tokio::time::sleep(Duration::from_millis(25)).await;
538 }
539 targets.send_replace(Vec::new());
540 targets
541 });
542 let executor = RacingStop::default();
543 let cause = "this agent does not offer high as a effort";
544
545 let outcome = loaded_controller()
546 .fail_subagent_start_worker(LATCH_RELAY_SESSION, cause, &executor, &control)
547 .await
548 .unwrap();
549
550 assert_eq!(outcome, ParkOutcome::Parked);
551 assert!(
552 executor
553 .purposes
554 .lock()
555 .unwrap()
556 .iter()
557 .any(|purpose| purpose == "stop Mjolnir worker daemon"),
558 "the child's worker was stopped: {:?}",
559 executor.purposes.lock().unwrap()
560 );
561 let stored = crate::database::load_state().unwrap();
562 let child = &stored.sessions[LATCH_RELAY_SESSION];
563 assert_eq!(child.state, SessionState::Error);
564 assert_eq!(child.last_error.as_deref(), Some(cause));
565 assert!(child.target.is_some(), "the failed child keeps its target");
566 assert!(stored.subagents.contains_key(LATCH_RELAY_SESSION));
567 assert!(control.session(LATCH_RELAY_SESSION).await.is_err());
568 let _targets = refresher.await.unwrap();
569 shutdown.shutdown().await.unwrap();
570 }
571
572 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
573 async fn a_child_with_work_in_flight_is_not_parked() {
574 if !isolated("a_child_with_work_in_flight_is_not_parked") {
575 return;
576 }
577 let _writer = crate::database::install_isolated_test_writer();
578 let root = tempfile::tempdir().unwrap();
579 register_child(root.path());
580 let channels = crate::session_manager::spawn_session_manager().unwrap();
581 channels
582 .targets
583 .send(vec![latch_relay_target(
584 root.path(),
585 None,
586 ReleaseSupport::Supported,
587 true,
588 )])
589 .unwrap();
590 channels
591 .control
592 .wait_for_session(LATCH_RELAY_SESSION, Duration::from_secs(10))
593 .await
594 .unwrap();
595 let executor = RacingStop::default();
596
597 let outcome = loaded_controller()
598 .park_subagent_worker(LATCH_RELAY_SESSION, &executor, &channels.control)
599 .await
600 .unwrap();
601
602 assert_eq!(outcome, ParkOutcome::Busy);
603 assert!(
604 executor.purposes.lock().unwrap().is_empty(),
605 "nothing stopped"
606 );
607 assert_eq!(
608 crate::database::load_session_state(LATCH_RELAY_SESSION).unwrap(),
609 Some(SessionState::Running)
610 );
611 channels.shutdown.shutdown().await.unwrap();
612 }
613
614 #[test]
615 fn only_a_full_target_is_rewritten_and_a_bare_one_reads_no_container_counts() {
616 let backend = targets::TargetLocator::LocalBare {
617 worker_root: "/tmp/workers/child".into(),
618 };
619 let unrelated =
620 explain_process_exhaustion(anyhow::anyhow!("the harness exited"), &backend, "child");
621 assert_eq!(format!("{unrelated:#}"), "the harness exited");
622
623 let full = explain_process_exhaustion(
624 anyhow::anyhow!("sh: 1: Cannot fork").context("start the parked sub-agent's worker"),
625 &backend,
626 "child",
627 );
628 let message = format!("{full:#}");
629 assert!(
630 message.starts_with("the target machine ran out of process slots"),
631 "{message}"
632 );
633 assert!(
634 message.contains("Close sub-agents you no longer need"),
635 "{message}"
636 );
637 assert!(
638 message.contains("Cannot fork"),
639 "the original error stays: {message}"
640 );
641 }
642}