Skip to main content

mj_controller/daemon/
subagent_park.rs

1//! The daemon's side of parking sub-agents (#1161): each park and unpark is a
2//! lifecycle operation of the child, so the per-session lifecycle map
3//! serializes it with the child's close, suspend, destroy and every other
4//! park or unpark. The work itself is in `crate::controller::subagent_park`.
5
6use super::*;
7
8/// How long a whole park may run: an idle reservation, one stop, one write.
9const PARK_TIMEOUT: Duration = Duration::from_secs(120);
10
11/// How long a whole unpark may run. It covers the harness start, which a
12/// worker restart allows 300 seconds for journal recovery alone.
13const UNPARK_TIMEOUT: Duration = Duration::from_secs(600);
14
15impl RuntimeState {
16    /// Park a sub-agent whose turn ended and whose parent was told. A child
17    /// that has work in flight or queued, or that is no longer running, is
18    /// left as it is. Any failure leaves the child running; the caller logs it.
19    pub async fn park_subagent(self: &Arc<Self>, child_session_id: String) -> Result<()> {
20        self.run_lifecycle(
21            child_session_id,
22            LifecycleKind::Park,
23            |state, session_id, _cancelled| async move {
24                // A park is short and not cancellable: a close that asks for
25                // the child waits for it instead, so it never finds a worker
26                // stopped under a record that still says running.
27                let never = AtomicBool::new(false);
28                let _recovery_reservation = blocking({
29                    let observer = state.recovery_observer.clone();
30                    let session_id = session_id.clone();
31                    move || reserve_recovery_or_cancel(&observer, &session_id, &never)
32                })
33                .await?;
34                if state.close_is_requested(&session_id) {
35                    return Ok(DaemonLifecycleResult::Done);
36                }
37                let controller = blocking(Controller::load).await?;
38                let executor = CancellableProcessExecutor::with_timeout(PARK_TIMEOUT);
39                let outcome = controller
40                    .park_subagent_worker(&session_id, &executor, &state.session_manager)
41                    .await?;
42                tracing::info!(%session_id, ?outcome, "sub-agent park finished");
43                Ok(DaemonLifecycleResult::Done)
44            },
45        )
46        .await?;
47        self.publish_revision();
48        Ok(())
49    }
50
51    /// Start a parked sub-agent's worker again so it can take its parent's
52    /// next prompt, and wait until the session manager holds it. A child that
53    /// is already running needs nothing. On failure the child stays parked.
54    pub async fn unpark_subagent(self: &Arc<Self>, child_session_id: String) -> Result<()> {
55        // A park still finishing would otherwise refuse this as a second
56        // lifecycle operation; waiting for it keeps a `send_input` that raced
57        // the park from failing.
58        self.wait_for_subagent_park(&child_session_id).await;
59        let parked = self
60            .session_state(&child_session_id)
61            .is_some_and(|state| state == SessionState::Parked);
62        if parked {
63            self.run_lifecycle(
64                child_session_id.clone(),
65                LifecycleKind::Unpark,
66                // A close of the child cancels this: the child then stays
67                // parked and the close settles it.
68                |state, session_id, cancelled| async move {
69                    let _recovery_reservation = blocking({
70                        let observer = state.recovery_observer.clone();
71                        let session_id = session_id.clone();
72                        let cancelled = cancelled.clone();
73                        move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
74                    })
75                    .await?;
76                    let controller = blocking(Controller::load).await?;
77                    let executor =
78                        CancellableProcessExecutor::new(cancelled).with_deadline(UNPARK_TIMEOUT);
79                    controller
80                        .unpark_subagent_worker(&session_id, &executor)
81                        .await?;
82                    Ok(DaemonLifecycleResult::Done)
83                },
84            )
85            .await?;
86            self.publish_revision();
87        }
88        Ok(())
89    }
90
91    /// Whether a park of this child has started and not finished.
92    pub fn subagent_park_running(&self, child_session_id: &str) -> bool {
93        self.lifecycle
94            .lock()
95            .unwrap_or_else(PoisonError::into_inner)
96            .get(child_session_id)
97            .is_some_and(|active| {
98                active.kind == LifecycleKind::Park && active.result.borrow().is_none()
99            })
100    }
101
102    /// Wait for a park of this child that is still running, if there is one.
103    async fn wait_for_subagent_park(self: &Arc<Self>, child_session_id: &str) {
104        let pending = self
105            .lifecycle
106            .lock()
107            .unwrap_or_else(PoisonError::into_inner)
108            .get(child_session_id)
109            .filter(|active| active.kind == LifecycleKind::Park)
110            .map(|active| active.result.clone());
111        if let Some(pending) = pending {
112            if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
113                tracing::debug!(
114                    session_id = %child_session_id,
115                    error = format!("{error:#}"),
116                    "the sub-agent park a restart waited for failed"
117                );
118            }
119            self.remove_completed_lifecycle(&pending);
120        }
121    }
122}