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