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
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
//! 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> {
self.stop_idle_subagent(child_session_id, None).await
}
/// Record a sub-agent whose first prompt was refused for good as failed
/// with `cause`, stopping its worker (I1-2). It runs as the child's park
/// operation: the same idle stop, serialized with its close and every
/// other lifecycle operation. A child that took work meanwhile is left
/// running. Failures are logged; the parent already reads the cause from
/// the failed startup whatever happens here.
pub(super) async fn fail_subagent_start(self: &Arc<Self>, child_session_id: &str, cause: &str) {
if !self
.owner()
.controller()
.state
.subagents
.contains_key(child_session_id)
{
return;
}
match self
.stop_idle_subagent(child_session_id.to_owned(), Some(cause.to_owned()))
.await
{
Ok(crate::controller::ParkOutcome::Parked) => {
tracing::info!(
session_id = child_session_id,
cause,
"sub-agent recorded as failed: its first prompt was refused"
);
}
Ok(outcome) => tracing::warn!(
session_id = child_session_id,
?outcome,
"a sub-agent whose first prompt was refused was left running"
),
Err(error) => tracing::warn!(
session_id = child_session_id,
error = format!("{error:#}"),
"could not record a sub-agent whose first prompt was refused as failed"
),
}
}
/// The park lifecycle: stop an idle child's worker and record it
/// `Parked`, or `Error` with `failure`.
async fn stop_idle_subagent(
self: &Arc<Self>,
child_session_id: String,
failure: Option<String>,
) -> Result<crate::controller::ParkOutcome> {
let result = self
.run_lifecycle(
child_session_id,
LifecycleKind::Park,
move |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 = match &failure {
None => {
controller
.park_subagent_worker(
&session_id,
&executor,
&state.session_manager,
)
.await?
}
Some(cause) => {
controller
.fail_subagent_start_worker(
&session_id,
cause,
&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);
}
}
}