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