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    const WAIT_PROMPT_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(1);
17
18    pub(crate) fn install_wait_prompt_backend(
19        &self,
20        backend: &Arc<crate::server_runtime::api::ApiBackend>,
21    ) {
22        assert!(
23            self.wait_prompt_backend
24                .set(Arc::downgrade(backend))
25                .is_ok(),
26            "delegation wait-prompt backend is installed once"
27        );
28    }
29
30    pub(crate) async fn ensure_parent_wait_prompt(&self, parent_session_id: &str) -> Result<()> {
31        if let Some(backend) = self
32            .wait_prompt_backend
33            .get()
34            .and_then(std::sync::Weak::upgrade)
35        {
36            match backend.ensure_parent_wait_prompt(parent_session_id).await {
37                Ok(()) => {
38                    self.wait_prompt_retries
39                        .lock()
40                        .unwrap_or_else(PoisonError::into_inner)
41                        .remove(parent_session_id);
42                }
43                Err(error) => {
44                    self.schedule_wait_prompt_retry(parent_session_id);
45                    return Err(error);
46                }
47            }
48        }
49        Ok(())
50    }
51
52    pub(crate) fn schedule_wait_prompt_retry(&self, parent_session_id: &str) {
53        self.wait_prompt_retries
54            .lock()
55            .unwrap_or_else(PoisonError::into_inner)
56            .insert(
57                parent_session_id.to_owned(),
58                std::time::Instant::now() + Self::WAIT_PROMPT_RETRY_DELAY,
59            );
60    }
61
62    pub(crate) fn take_due_wait_prompt_retries(
63        &self,
64        now: std::time::Instant,
65        limit: usize,
66    ) -> Vec<String> {
67        let mut retries = self
68            .wait_prompt_retries
69            .lock()
70            .unwrap_or_else(PoisonError::into_inner);
71        let due = retries
72            .iter()
73            .filter(|(_, due)| **due <= now)
74            .take(limit)
75            .map(|(parent, _)| parent.clone())
76            .collect::<Vec<_>>();
77        for parent in &due {
78            retries.remove(parent);
79        }
80        due
81    }
82
83    /// Park a sub-agent whose turn ended and whose parent was told. A child
84    /// that has work in flight or queued, or that is no longer running, is
85    /// left as it is. Any failure leaves the child running; the caller logs it.
86    pub async fn park_subagent(
87        self: &Arc<Self>,
88        child_session_id: String,
89    ) -> Result<crate::controller::ParkOutcome> {
90        self.stop_idle_subagent(child_session_id, None).await
91    }
92
93    /// Record a sub-agent whose first prompt was refused for good as failed
94    /// with `cause`, stopping its worker (I1-2). It runs as the child's park
95    /// operation: the same idle stop, serialized with its close and every
96    /// other lifecycle operation. A child that took work meanwhile is left
97    /// running. Failures are logged; the parent already reads the cause from
98    /// the failed startup whatever happens here.
99    pub(super) async fn fail_subagent_start(self: &Arc<Self>, child_session_id: &str, cause: &str) {
100        let parent_session_id = self
101            .owner()
102            .controller()
103            .state
104            .subagents
105            .get(child_session_id)
106            .map(|relation| relation.parent_session_id.clone());
107        if !self
108            .owner()
109            .controller()
110            .state
111            .subagents
112            .contains_key(child_session_id)
113        {
114            return;
115        }
116        match self
117            .stop_idle_subagent(child_session_id.to_owned(), Some(cause.to_owned()))
118            .await
119        {
120            Ok(crate::controller::ParkOutcome::Parked) => {
121                tracing::info!(
122                    session_id = child_session_id,
123                    cause,
124                    "sub-agent recorded as failed: its first prompt was refused"
125                );
126            }
127            Ok(outcome) => tracing::warn!(
128                session_id = child_session_id,
129                ?outcome,
130                "a sub-agent whose first prompt was refused was left running"
131            ),
132            Err(error) => tracing::warn!(
133                session_id = child_session_id,
134                error = format!("{error:#}"),
135                "could not record a sub-agent whose first prompt was refused as failed"
136            ),
137        }
138        if let Some(parent_session_id) = parent_session_id
139            && let Err(error) = self.ensure_parent_wait_prompt(&parent_session_id).await
140        {
141            tracing::warn!(
142                child_session_id,
143                parent_session_id,
144                error = format!("{error:#}"),
145                "could not reconcile the parent's sub-agent wait prompt after startup failure"
146            );
147        }
148    }
149
150    /// The park lifecycle: stop an idle child's worker and record it
151    /// `Parked`, or `Error` with `failure`.
152    async fn stop_idle_subagent(
153        self: &Arc<Self>,
154        child_session_id: String,
155        failure: Option<String>,
156    ) -> Result<crate::controller::ParkOutcome> {
157        let result = self
158            .run_lifecycle(
159                child_session_id,
160                LifecycleKind::Park,
161                move |state, session_id, _cancelled| async move {
162                    // A park is short and not cancellable: a close that asks for
163                    // the child waits for it instead, so it never finds a worker
164                    // stopped under a record that still says running.
165                    let never = AtomicBool::new(false);
166                    let _recovery_reservation = blocking({
167                        let observer = state.recovery_observer.clone();
168                        let session_id = session_id.clone();
169                        move || reserve_recovery_or_cancel(&observer, &session_id, &never)
170                    })
171                    .await?;
172                    if state.close_is_requested(&session_id) {
173                        return Ok(DaemonLifecycleResult::Park(
174                            crate::controller::ParkOutcome::NotRunning,
175                        ));
176                    }
177                    let controller = blocking(Controller::load).await?;
178                    let executor = CancellableProcessExecutor::with_timeout(PARK_TIMEOUT);
179                    let outcome = match &failure {
180                        None => {
181                            controller
182                                .park_subagent_worker(
183                                    &session_id,
184                                    &executor,
185                                    &state.session_manager,
186                                )
187                                .await?
188                        }
189                        Some(cause) => {
190                            controller
191                                .fail_subagent_start_worker(
192                                    &session_id,
193                                    cause,
194                                    &executor,
195                                    &state.session_manager,
196                                )
197                                .await?
198                        }
199                    };
200                    tracing::info!(%session_id, ?outcome, "sub-agent park finished");
201                    Ok(DaemonLifecycleResult::Park(outcome))
202                },
203            )
204            .await?;
205        self.publish_revision();
206        match result {
207            DaemonLifecycleResult::Park(outcome) => Ok(outcome),
208            _ => unreachable!("a park returns its worker's reservation outcome"),
209        }
210    }
211
212    /// Start a parked sub-agent's worker again so it can take its parent's
213    /// next prompt, and wait until the session manager holds it. A child that
214    /// is already running needs nothing. On failure the child stays parked.
215    pub async fn unpark_subagent(self: &Arc<Self>, child_session_id: String) -> Result<()> {
216        // A park still finishing would otherwise refuse this as a second
217        // lifecycle operation; waiting for it keeps a `send_input` that raced
218        // the park from failing.
219        self.wait_for_subagent_park(&child_session_id).await;
220        let parked = self
221            .session_state(&child_session_id)
222            .is_some_and(|state| state == SessionState::Parked);
223        if parked {
224            self.run_lifecycle(
225                child_session_id.clone(),
226                LifecycleKind::Unpark,
227                // A close of the child cancels this: the child then stays
228                // parked and the close settles it.
229                |state, session_id, cancelled| async move {
230                    let _recovery_reservation = blocking({
231                        let observer = state.recovery_observer.clone();
232                        let session_id = session_id.clone();
233                        let cancelled = cancelled.clone();
234                        move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
235                    })
236                    .await?;
237                    let controller = blocking(Controller::load).await?;
238                    let executor =
239                        CancellableProcessExecutor::new(cancelled).with_deadline(UNPARK_TIMEOUT);
240                    controller
241                        .unpark_subagent_worker(&session_id, &executor)
242                        .await?;
243                    Ok(DaemonLifecycleResult::Done)
244                },
245            )
246            .await?;
247            self.publish_revision();
248        }
249        Ok(())
250    }
251
252    /// Whether a park of this child has started and not finished.
253    pub fn subagent_park_running(&self, child_session_id: &str) -> bool {
254        self.owner()
255            .lifecycle
256            .get(child_session_id)
257            .is_some_and(|active| active.kind == LifecycleKind::Park && active.is_running())
258    }
259
260    /// Wait for a park of this child that is still running, if there is one.
261    async fn wait_for_subagent_park(self: &Arc<Self>, child_session_id: &str) {
262        let pending = self
263            .owner()
264            .lifecycle
265            .get(child_session_id)
266            .filter(|active| active.kind == LifecycleKind::Park)
267            .map(|active| active.result.clone());
268        if let Some(pending) = pending {
269            if let Err(error) = Self::wait_lifecycle_result(pending.clone()).await {
270                tracing::debug!(
271                    session_id = %child_session_id,
272                    error = format!("{error:#}"),
273                    "the sub-agent park a restart waited for failed"
274                );
275            }
276            self.remove_completed_lifecycle(&pending);
277        }
278    }
279}