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