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(
20        self: &Arc<Self>,
21        child_session_id: String,
22    ) -> Result<crate::controller::ParkOutcome> {
23        self.stop_idle_subagent(child_session_id, None).await
24    }
25
26    /// Record a sub-agent whose first prompt was refused for good as failed
27    /// with `cause`, stopping its worker (I1-2). It runs as the child's park
28    /// operation: the same idle stop, serialized with its close and every
29    /// other lifecycle operation. A child that took work meanwhile is left
30    /// running. Failures are logged; the parent already reads the cause from
31    /// the failed startup whatever happens here.
32    pub(super) async fn fail_subagent_start(self: &Arc<Self>, child_session_id: &str, cause: &str) {
33        if !self
34            .owner()
35            .controller()
36            .state
37            .subagents
38            .contains_key(child_session_id)
39        {
40            return;
41        }
42        match self
43            .stop_idle_subagent(child_session_id.to_owned(), Some(cause.to_owned()))
44            .await
45        {
46            Ok(crate::controller::ParkOutcome::Parked) => {
47                tracing::info!(
48                    session_id = child_session_id,
49                    cause,
50                    "sub-agent recorded as failed: its first prompt was refused"
51                );
52            }
53            Ok(outcome) => tracing::warn!(
54                session_id = child_session_id,
55                ?outcome,
56                "a sub-agent whose first prompt was refused was left running"
57            ),
58            Err(error) => tracing::warn!(
59                session_id = child_session_id,
60                error = format!("{error:#}"),
61                "could not record a sub-agent whose first prompt was refused as failed"
62            ),
63        }
64    }
65
66    /// The park lifecycle: stop an idle child's worker and record it
67    /// `Parked`, or `Error` with `failure`.
68    async fn stop_idle_subagent(
69        self: &Arc<Self>,
70        child_session_id: String,
71        failure: Option<String>,
72    ) -> Result<crate::controller::ParkOutcome> {
73        let result = self
74            .run_lifecycle(
75                child_session_id,
76                LifecycleKind::Park,
77                move |state, session_id, _cancelled| async move {
78                    // A park is short and not cancellable: a close that asks for
79                    // the child waits for it instead, so it never finds a worker
80                    // stopped under a record that still says running.
81                    let never = AtomicBool::new(false);
82                    let _recovery_reservation = blocking({
83                        let observer = state.recovery_observer.clone();
84                        let session_id = session_id.clone();
85                        move || reserve_recovery_or_cancel(&observer, &session_id, &never)
86                    })
87                    .await?;
88                    if state.close_is_requested(&session_id) {
89                        return Ok(DaemonLifecycleResult::Park(
90                            crate::controller::ParkOutcome::NotRunning,
91                        ));
92                    }
93                    let controller = blocking(Controller::load).await?;
94                    let executor = CancellableProcessExecutor::with_timeout(PARK_TIMEOUT);
95                    let outcome = match &failure {
96                        None => {
97                            controller
98                                .park_subagent_worker(
99                                    &session_id,
100                                    &executor,
101                                    &state.session_manager,
102                                )
103                                .await?
104                        }
105                        Some(cause) => {
106                            controller
107                                .fail_subagent_start_worker(
108                                    &session_id,
109                                    cause,
110                                    &executor,
111                                    &state.session_manager,
112                                )
113                                .await?
114                        }
115                    };
116                    tracing::info!(%session_id, ?outcome, "sub-agent park finished");
117                    Ok(DaemonLifecycleResult::Park(outcome))
118                },
119            )
120            .await?;
121        self.publish_revision();
122        match result {
123            DaemonLifecycleResult::Park(outcome) => Ok(outcome),
124            _ => unreachable!("a park returns its worker's reservation outcome"),
125        }
126    }
127
128    /// Start a parked sub-agent's worker again so it can take its parent's
129    /// next prompt, and wait until the session manager holds it. A child that
130    /// is already running needs nothing. On failure the child stays parked.
131    pub async fn unpark_subagent(self: &Arc<Self>, child_session_id: String) -> Result<()> {
132        // A park still finishing would otherwise refuse this as a second
133        // lifecycle operation; waiting for it keeps a `send_input` that raced
134        // the park from failing.
135        self.wait_for_subagent_park(&child_session_id).await;
136        let parked = self
137            .session_state(&child_session_id)
138            .is_some_and(|state| state == SessionState::Parked);
139        if parked {
140            self.run_lifecycle(
141                child_session_id.clone(),
142                LifecycleKind::Unpark,
143                // A close of the child cancels this: the child then stays
144                // parked and the close settles it.
145                |state, session_id, cancelled| async move {
146                    let _recovery_reservation = blocking({
147                        let observer = state.recovery_observer.clone();
148                        let session_id = session_id.clone();
149                        let cancelled = cancelled.clone();
150                        move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
151                    })
152                    .await?;
153                    let controller = blocking(Controller::load).await?;
154                    let executor =
155                        CancellableProcessExecutor::new(cancelled).with_deadline(UNPARK_TIMEOUT);
156                    controller
157                        .unpark_subagent_worker(&session_id, &executor)
158                        .await?;
159                    Ok(DaemonLifecycleResult::Done)
160                },
161            )
162            .await?;
163            self.publish_revision();
164        }
165        Ok(())
166    }
167
168    /// Whether a park of this child has started and not finished.
169    pub fn subagent_park_running(&self, child_session_id: &str) -> bool {
170        self.owner()
171            .lifecycle
172            .get(child_session_id)
173            .is_some_and(|active| active.kind == LifecycleKind::Park && active.is_running())
174    }
175
176    /// Wait for a park of this child that is still running, if there is one.
177    async fn wait_for_subagent_park(self: &Arc<Self>, child_session_id: &str) {
178        let pending = self
179            .owner()
180            .lifecycle
181            .get(child_session_id)
182            .filter(|active| active.kind == LifecycleKind::Park)
183            .map(|active| active.result.clone());
184        if let Some(pending) = pending {
185            if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
186                tracing::debug!(
187                    session_id = %child_session_id,
188                    error = format!("{error:#}"),
189                    "the sub-agent park a restart waited for failed"
190                );
191            }
192            self.remove_completed_lifecycle(&pending);
193        }
194    }
195}