Skip to main content

mj_controller/controller/
subagent_park.rs

1//! Parking a sub-agent whose turn has ended, and starting it again.
2//!
3//! A sub-agent child runs inside its parent's target, and in a container every
4//! live child's harness holds hundreds of threads against the container's pids
5//! limit (#1161). A child that has finished its turn, and whose parent has
6//! been told so, therefore gives its processes back: its worker is stopped and
7//! its record says [`SessionState::Parked`]. Everything else stays: the
8//! record, the relation to the parent, the borrowed target locator, and the
9//! worker root with its relay journal and the harness's native session id. A
10//! parked child is started again in place, with the same sequence a worker
11//! restart uses, when its parent gives it more input.
12//!
13//! The daemon runs both as lifecycle operations of the child, so they are
14//! serialized with its close, suspend, destroy and each other.
15
16use 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
31/// How long a park waits for the child's session actor to exist.
32const PARK_ACTOR_TIMEOUT: Duration = Duration::from_secs(5);
33
34/// How long a park waits, after recording `Parked`, for the session manager
35/// to drop the child. The manager rereads the store twice a second.
36const PARK_RELEASE_TIMEOUT: Duration = Duration::from_secs(15);
37
38/// What an unpark tells the operator at each step of the restart it runs.
39const 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/// How one park attempt ended. Only [`ParkOutcome::Parked`] changed anything.
49#[derive(Debug, Clone, Copy, PartialEq, Eq)]
50pub enum ParkOutcome {
51    /// The worker was stopped and the record says `Parked`.
52    Parked,
53    /// The child had work in flight or queued, such as a prompt its parent
54    /// sent after the turn ended, so it was left running.
55    Busy,
56    /// The child is not a running sub-agent any more (it is closing, already
57    /// parked, or gone), so there was nothing to park.
58    NotRunning,
59}
60
61/// An executor for stopping what a failed unpark started. The unpark's own
62/// executor may be the reason it failed, cancelled by a close of the child.
63fn cleanup_executor() -> targets::CancellableProcessExecutor {
64    targets::CancellableProcessExecutor::with_timeout(Duration::from_secs(30))
65}
66
67/// Say plainly that a sub-agent could not start because its target ran out of
68/// process slots, when `error` shows that, with the container's pid counts
69/// when they can be read; return any other error unchanged.
70///
71/// The read is best effort, bounded by [`targets::PIDS_USAGE_READ_TIMEOUT`],
72/// and itself fails in a container that is completely full, which the
73/// message then says.
74pub(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    /// Stop a running sub-agent's worker while it is idle, and record it as
92    /// parked.
93    ///
94    /// The worker is reserved the way an idle worker upgrade reserves it: the
95    /// actor's connection is leased only while the worker reports nothing
96    /// running or queued, and the worker holds an idle barrier until it is
97    /// stopped. A prompt that reaches the actor meanwhile waits behind the
98    /// lease. The lease is kept until the session manager has dropped the
99    /// child, which it does once the store says `Parked`, so such a prompt is
100    /// rejected as undelivered rather than sent to a stopped worker; the
101    /// caller that sent it can start the child again and resend it.
102    ///
103    /// Any failure before the record changes leaves the child running: the
104    /// lease is dropped and the actor reconnects, restarting the worker if the
105    /// stop got that far.
106    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    /// Stop a sub-agent whose first prompt can never run and record it as
117    /// failed with `cause`, so every surface shows the failure instead of an
118    /// idle session holding a worker (I1-2). It is stopped exactly as a park
119    /// stops a child, behind the same idle reservation, and keeps its record,
120    /// relation and target; its close settles without a checkpoint (see
121    /// [`super::has_nothing_to_checkpoint`]).
122    ///
123    /// A child that took work meanwhile is left running and
124    /// [`ParkOutcome::Busy`] is returned. [`ParkOutcome::Parked`] means the
125    /// worker was stopped and the record now says `Error`.
126    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    /// The stop behind a park and a failed start: the record becomes
138    /// `Parked`, or `Error` with `failure` as its cause.
139    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    /// Start a parked sub-agent's worker again in place and record it as
196    /// running.
197    ///
198    /// The worker binary and launch configuration are refreshed first when
199    /// this controller would now install different ones, then the restart
200    /// sequence runs: start the worker on its existing root, connect with the
201    /// long restart timeout, and wait until the harness has loaded its native
202    /// session and is idle. The container's start admission is held around
203    /// the harness start, as it is for a child's first start.
204    ///
205    /// A failure stops whatever was started and leaves the record `Parked`,
206    /// so the parent can try again. Nothing here connects a session actor:
207    /// the caller runs this while the daemon keeps the manager off the child,
208    /// and the manager attaches once the record says `Running`.
209    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            // Held across the harness start, which is the part that does not
244            // survive a crowd of children starting in one container.
245            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        // The manager opens its own connection once the record says running.
279        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    /// Run the named test alone, with a store of its own. Returns whether
316    /// this process is that run.
317    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    /// Stands in for the target: every command succeeds, and the first one,
330    /// which is the park's stop, also sends the child a prompt through its
331    /// actor, the way a parent's `send_input` can race a park.
332    #[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    /// A parent, and the stand-in relay's session registered as its child on
368    /// a bare target under `root`, with a report it handed back.
369    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        // What the daemon's target refresher does: once the store says the
443        // child is parked, the session manager no longer holds it.
444        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        // The prompt that arrived during the park never reached the stopped
471        // worker, and its sender is told so for certain, so it can start the
472        // child again and send it once more.
473        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        // Everything but the worker stays.
486        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    /// I1-2: a child whose first prompt is refused for good has its worker
505    /// stopped and is recorded as failed with the cause, keeping its record,
506    /// relation and target, instead of staying a live idle session.
507    #[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}