refactor(agent): fold parent notifications into subagent state

`ParentNotificationHub` kept a second `Mutex` and `watch` channel holding
a copy of each child's terminal result -- data `SubAgent.status` already
owns as `SubAgentStatus::Finished`, and which is never evicted, since
nothing removes entries from `SupervisorState.agents`.

Two of the three bugs fixed in the previous commit were ordering bugs in
the coupling between those two structures: suppress-vs-commit in
`begin_shutdown`, and register-vs-publish in `spawn_inner`. Both were
fixed by ordering the steps correctly. Keeping the registration beside
the status it is delivered with makes that whole class unrepresentable
instead:

- Registration is now a field on the `SubAgent` literal `spawn_inner`
  already builds, under the lock that publishes it. There is no window
  between publishing an agent and registering its notification.
- Suppression on shutdown happens inside the critical section that
  decides the shutdown, after the status transition commits, so a
  rejected shutdown cannot discard a result the parent is owed.
- `next_parent_notification_batch` scans agents for a live registration
  whose status is `Finished`, and ignores `Closing`/`Closed` outright --
  so a shutdown racing delivery can no longer park the parent on a result
  that will never arrive, even if suppression were missed.

`spawn_result_monitor` no longer takes the hub; it bumps a single
`watch` counter after committing the status it already commits. Batch
order was the queue's insertion order, so `SubAgent` carries a
`spawn_seq` to keep delivery oldest-first.

Tests move from exercising the hub directly to the supervisor API, and
cover spawn-order batching and the shutdown-races-delivery case that the
old shape could not express.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-07-25 23:44:20 -04:00
parent d669a2d55c
commit f293e3de18
No known key found for this signature in database

View file

@ -89,16 +89,27 @@ pub enum SubAgentStatus {
const SUBAGENT_SHUTDOWN_GRACE: Duration = Duration::from_secs(5);
struct SubAgent {
status: watch::Sender<SubAgentStatus>,
cleanup_done: watch::Sender<bool>,
cleanup_started: bool,
monitor_task: Option<JoinHandle<()>>,
event_forwarder: Option<JoinHandle<()>>,
cleanup_task: Option<JoinHandle<()>>,
child_abort_handle: AbortHandle,
followup_queue: Arc<Mutex<VecDeque<String>>>,
cancel_token: CancellationToken,
depth: usize,
status: watch::Sender<SubAgentStatus>,
cleanup_done: watch::Sender<bool>,
cleanup_started: bool,
monitor_task: Option<JoinHandle<()>>,
event_forwarder: Option<JoinHandle<()>>,
cleanup_task: Option<JoinHandle<()>>,
child_abort_handle: AbortHandle,
followup_queue: Arc<Mutex<VecDeque<String>>>,
cancel_token: CancellationToken,
depth: usize,
/// Task description, set when the parent should receive this child's
/// terminal result automatically. Cleared once the result is delivered,
/// the parent retrieves it explicitly, or the agent is shut down.
///
/// Keeping this beside the status it is delivered with means a
/// notification cannot be registered before -- or suppressed after -- the
/// state it describes: there is only one lock and one ordering.
parent_notification: Option<String>,
/// Spawn order, so a batch is delivered oldest-first rather than in
/// whatever order the map happens to iterate.
spawn_seq: u64,
}
impl Drop for SubAgent {
@ -119,114 +130,8 @@ impl Drop for SubAgent {
#[derive(Default)]
struct SupervisorState {
agents: HashMap<String, SubAgent>,
}
#[derive(Default)]
struct ParentNotificationState {
pending: HashMap<String, String>,
ready: VecDeque<SubAgentParentNotification>,
}
struct ParentNotificationHub {
state: Mutex<ParentNotificationState>,
changed: watch::Sender<u64>,
}
impl ParentNotificationHub {
fn new() -> Self {
let (changed, _) = watch::channel(0);
Self {
state: Mutex::new(ParentNotificationState::default()),
changed,
}
}
fn register(&self, agent_id: String, description: String) {
self.state
.lock()
.expect("parent notification lock poisoned")
.pending
.insert(agent_id, description);
self.signal();
}
fn complete(&self, agent_id: &str, result: Result<SubAgentResult, Error>) {
{
let mut state = self
.state
.lock()
.expect("parent notification lock poisoned");
let Some(description) = state.pending.remove(agent_id) else {
return;
};
state.ready.push_back(SubAgentParentNotification {
agent_id: agent_id.to_string(),
description,
result,
});
}
self.signal();
}
fn suppress(&self, agent_id: &str) {
let changed = {
let mut state = self
.state
.lock()
.expect("parent notification lock poisoned");
let removed_pending = state.pending.remove(agent_id).is_some();
let ready_len = state.ready.len();
state
.ready
.retain(|notification| notification.agent_id != agent_id);
removed_pending || state.ready.len() != ready_len
};
if changed {
self.signal();
}
}
async fn next_batch(
&self,
cancel: &CancellationToken,
) -> Result<Option<Vec<SubAgentParentNotification>>, Error> {
let mut changed = self.changed.subscribe();
loop {
{
let mut state = self
.state
.lock()
.expect("parent notification lock poisoned");
if !state.ready.is_empty() {
return Ok(Some(state.ready.drain(..).collect()));
}
if state.pending.is_empty() {
return Ok(None);
}
}
tokio::select! {
biased;
() = cancel.cancelled() => {
return Err(Error::Interrupted(InterruptReason::Cancelled));
}
observed = changed.changed() => {
observed.map_err(|_| {
Error::InvalidState(
"Background-agent notification observer closed unexpectedly".to_string(),
)
})?;
}
}
}
}
fn signal(&self) {
self.changed.send_modify(|generation| {
*generation = generation.wrapping_add(1);
});
}
agents: HashMap<String, SubAgent>,
next_spawn_seq: u64,
}
struct ShutdownWork {
@ -268,11 +173,20 @@ impl Drop for CleanupDoneGuard {
}
}
/// Wake anything parked in
/// [`SubAgentSupervisor::next_parent_notification_batch`] so it can re-evaluate
/// which children are deliverable.
fn signal_notifications(changed: &watch::Sender<u64>) {
changed.send_modify(|generation| {
*generation = generation.wrapping_add(1);
});
}
fn spawn_result_monitor(
child_task: JoinHandle<Result<SubAgentResult, Error>>,
status: watch::Sender<SubAgentStatus>,
event_callback: Arc<RwLock<Option<SubAgentEventCallback>>>,
parent_notifications: Arc<ParentNotificationHub>,
notifications_changed: Arc<watch::Sender<u64>>,
agent_id: String,
depth: usize,
) -> JoinHandle<()> {
@ -294,6 +208,8 @@ fn spawn_result_monitor(
if !committed {
return;
}
// The status this agent will be delivered with is now committed.
signal_notifications(&notifications_changed);
let event = match &task_result {
Ok(result) => AgentEvent::SubAgentCompleted {
@ -315,7 +231,6 @@ fn spawn_result_monitor(
if let Some(callback) = callback {
callback(SubAgentCallbackEvent::Lifecycle(event));
}
parent_notifications.complete(&agent_id, task_result);
})
}
@ -326,10 +241,10 @@ fn spawn_result_monitor(
/// happen after the guard has been released.
#[derive(Clone)]
pub struct SubAgentSupervisor {
state: Arc<Mutex<SupervisorState>>,
max_depth: usize,
event_callback: Arc<RwLock<Option<SubAgentEventCallback>>>,
parent_notifications: Arc<ParentNotificationHub>,
state: Arc<Mutex<SupervisorState>>,
max_depth: usize,
event_callback: Arc<RwLock<Option<SubAgentEventCallback>>>,
notifications_changed: Arc<watch::Sender<u64>>,
}
impl SubAgentSupervisor {
@ -339,7 +254,7 @@ impl SubAgentSupervisor {
state: Arc::new(Mutex::new(SupervisorState::default())),
max_depth,
event_callback: Arc::new(RwLock::new(None)),
parent_notifications: Arc::new(ParentNotificationHub::new()),
notifications_changed: Arc::new(watch::channel(0).0),
}
}
@ -474,23 +389,15 @@ impl SubAgentSupervisor {
child_task,
status.clone(),
Arc::clone(&self.event_callback),
Arc::clone(&self.parent_notifications),
Arc::clone(&self.notifications_changed),
agent_id.clone(),
child_depth,
);
// Register before the agent becomes discoverable in `state.agents`.
// Once it is, a concurrent `shutdown_all` can suppress and close it;
// registering afterwards would leave a pending entry that the monitor
// never completes (it early-returns for a non-Running agent), and
// `next_batch` would then never report the queue as drained. Nothing
// can complete this registration before `start_tx.send(())` below.
if let Some(description) = parent_notification_description {
self.parent_notifications
.register(agent_id.clone(), description);
}
{
let mut state = self.state.lock().expect("subagent state lock poisoned");
let spawn_seq = state.next_spawn_seq;
state.next_spawn_seq = state.next_spawn_seq.saturating_add(1);
state.agents.insert(agent_id.clone(), SubAgent {
status,
cleanup_done,
@ -502,8 +409,11 @@ impl SubAgentSupervisor {
followup_queue,
cancel_token,
depth: child_depth,
parent_notification: parent_notification_description,
spawn_seq,
});
}
signal_notifications(&self.notifications_changed);
self.emit_event(AgentEvent::SubAgentSpawned {
agent_id: agent_id.clone(),
@ -587,11 +497,20 @@ impl SubAgentSupervisor {
}
}
/// Stop automatic delivery for an agent whose result the parent explicitly
/// retrieved. Removes a result that may already have raced into the ready
/// queue.
/// Stop automatic delivery for an agent whose result the parent retrieved
/// explicitly.
pub(crate) fn suppress_parent_notification(&self, agent_id: &str) {
self.parent_notifications.suppress(agent_id);
let cleared = {
let mut state = self.state.lock().expect("subagent state lock poisoned");
state
.agents
.get_mut(agent_id)
.and_then(|agent| agent.parent_notification.take())
.is_some()
};
if cleared {
signal_notifications(&self.notifications_changed);
}
}
/// Wait until all currently-ready background results can be delivered in
@ -616,7 +535,68 @@ impl SubAgentSupervisor {
&self,
cancel: &CancellationToken,
) -> Result<Option<Vec<SubAgentParentNotification>>, Error> {
self.parent_notifications.next_batch(cancel).await
let mut changed = self.notifications_changed.subscribe();
loop {
{
let mut state = self.state.lock().expect("subagent state lock poisoned");
let mut ready = Vec::new();
let mut awaiting_result = false;
for (agent_id, agent) in &state.agents {
let Some(description) = agent.parent_notification.as_ref() else {
continue;
};
let finished = match &*agent.status.borrow() {
SubAgentStatus::Finished(result) => Some(result.clone()),
SubAgentStatus::Running => {
awaiting_result = true;
None
}
// Being torn down, so no result is coming. Ignoring
// these is what keeps a shutdown that races delivery
// from parking the parent forever.
SubAgentStatus::Closing | SubAgentStatus::Closed => None,
};
if let Some(result) = finished {
ready.push((agent.spawn_seq, SubAgentParentNotification {
agent_id: agent_id.clone(),
description: description.clone(),
result,
}));
}
}
if !ready.is_empty() {
ready.sort_by_key(|(spawn_seq, _)| *spawn_seq);
let batch: Vec<_> = ready
.into_iter()
.map(|(_, notification)| notification)
.collect();
for notification in &batch {
if let Some(agent) = state.agents.get_mut(&notification.agent_id) {
agent.parent_notification = None;
}
}
return Ok(Some(batch));
}
if !awaiting_result {
return Ok(None);
}
}
tokio::select! {
biased;
() = cancel.cancelled() => {
return Err(Error::Interrupted(InterruptReason::Cancelled));
}
observed = changed.changed() => {
observed.map_err(|_| {
Error::InvalidState(
"Background-agent notification observer closed unexpectedly".to_string(),
)
})?;
}
}
}
}
#[cfg(test)]
@ -666,6 +646,11 @@ impl SubAgentSupervisor {
}
};
// Shutdown is committed, so this child's result will never reach the
// parent. The early returns above leave the notification intact, so a
// rejected shutdown cannot discard a result the parent is owed.
agent.parent_notification = None;
if agent.cleanup_started {
return Ok(ShutdownDisposition::Follow(agent.cleanup_done.subscribe()));
}
@ -772,10 +757,7 @@ impl SubAgentSupervisor {
async fn ensure_closed(&self, agent_id: &str) -> Result<(), Error> {
let disposition = self.begin_shutdown(agent_id, false)?;
// Only once shutdown is committed. Suppressing before `begin_shutdown`
// would also discard the result of an agent that had already finished,
// which rejects the shutdown but had a delivery pending.
self.parent_notifications.suppress(agent_id);
signal_notifications(&self.notifications_changed);
let cleanup_done = match disposition {
ShutdownDisposition::Lead(work) => self.spawn_shutdown(work),
ShutdownDisposition::Follow(cleanup_done) => cleanup_done,
@ -788,7 +770,7 @@ impl SubAgentSupervisor {
/// Strict user-facing close: only a currently running child may be closed.
pub async fn close_agent(&self, agent_id: &str) -> Result<(), Error> {
let disposition = self.begin_shutdown(agent_id, true)?;
self.parent_notifications.suppress(agent_id);
signal_notifications(&self.notifications_changed);
let cleanup_done = match disposition {
ShutdownDisposition::Lead(work) => self.spawn_shutdown(work),
ShutdownDisposition::Follow(_) | ShutdownDisposition::Done => {
@ -857,7 +839,7 @@ impl SubAgentSupervisor {
child_task,
status.clone(),
Arc::clone(&self.event_callback),
Arc::clone(&self.parent_notifications),
Arc::clone(&self.notifications_changed),
agent_id.clone(),
depth,
);
@ -873,6 +855,8 @@ impl SubAgentSupervisor {
event_forwarder,
cleanup_task: None,
child_abort_handle,
parent_notification: None,
spawn_seq: 0,
followup_queue: Arc::new(Mutex::new(VecDeque::new())),
cancel_token,
depth,
@ -1072,67 +1056,155 @@ mod tests {
assert!(manager.is_empty());
}
#[tokio::test]
async fn parent_notifications_are_exactly_once_and_xml_escaped() {
let hub = ParentNotificationHub::new();
hub.register("agent<&".to_string(), "Review <core> & tests".to_string());
let result = Ok(SubAgentResult {
output: "done <safely> & \"verified\"".to_string(),
success: true,
turns_used: 2,
});
hub.complete("agent<&", result.clone());
hub.complete("agent<&", result);
#[test]
fn parent_notification_envelope_escapes_xml() {
let envelope = format_parent_notification_batch(&[SubAgentParentNotification {
agent_id: "agent<&".to_string(),
description: "Review <core> & tests".to_string(),
result: Ok(SubAgentResult {
output: "done <safely> & \"verified\"".to_string(),
success: true,
turns_used: 2,
}),
}]);
let notifications = hub
.next_batch(&CancellationToken::new())
.await
.unwrap()
.unwrap();
assert_eq!(notifications.len(), 1);
let envelope = format_parent_notification_batch(&notifications);
assert!(envelope.contains("<status>completed</status>"));
assert!(envelope.contains("<task-id>agent&lt;&amp;</task-id>"));
assert!(envelope.contains("<description>Review &lt;core&gt; &amp; tests</description>"));
assert!(
envelope.contains("<result>done &lt;safely&gt; &amp; &quot;verified&quot;</result>")
);
assert!(
hub.next_batch(&CancellationToken::new())
.await
.unwrap()
.is_none()
);
}
#[tokio::test]
async fn suppress_removes_pending_and_ready_parent_notifications() {
let hub = ParentNotificationHub::new();
hub.register("pending".to_string(), "Pending".to_string());
hub.suppress("pending");
async fn a_finished_agent_is_delivered_to_the_parent_exactly_once() {
let supervisor = SubAgentSupervisor::new(3);
let child = make_session(vec![text_response("child result")]).await;
let agent_id = supervisor
.spawn_with_parent_notification(
child,
"task".to_string(),
"Inspect the module".to_string(),
0,
)
.unwrap();
supervisor
.wait_with_cancel(&agent_id, &CancellationToken::new())
.await
.unwrap();
let batch = supervisor
.next_parent_notification_batch(&CancellationToken::new())
.await
.unwrap()
.expect("the finished child must be delivered");
assert_eq!(batch.len(), 1);
assert_eq!(batch[0].agent_id, agent_id);
assert_eq!(batch[0].description, "Inspect the module");
// The status stays `Finished`, so re-delivery is prevented by clearing
// the registration rather than by consuming the result.
assert!(
hub.next_batch(&CancellationToken::new())
supervisor
.next_parent_notification_batch(&CancellationToken::new())
.await
.unwrap()
.is_none()
);
hub.register("ready".to_string(), "Ready".to_string());
hub.complete(
"ready",
Ok(SubAgentResult {
output: "done".to_string(),
success: true,
turns_used: 1,
}),
);
hub.suppress("ready");
supervisor.shutdown_all().await;
}
#[tokio::test]
async fn batches_are_delivered_in_spawn_order() {
let supervisor = SubAgentSupervisor::new(3);
let mut ids = Vec::new();
for index in 0..3 {
let child = make_session(vec![text_response("done")]).await;
ids.push(
supervisor
.spawn_with_parent_notification(
child,
format!("task {index}"),
format!("Task {index}"),
0,
)
.unwrap(),
);
}
for id in &ids {
supervisor
.wait_with_cancel(id, &CancellationToken::new())
.await
.unwrap();
}
let batch = supervisor
.next_parent_notification_batch(&CancellationToken::new())
.await
.unwrap()
.expect("all three children must be delivered together");
let delivered: Vec<_> = batch.iter().map(|n| n.agent_id.clone()).collect();
assert_eq!(delivered, ids);
supervisor.shutdown_all().await;
}
#[tokio::test]
async fn suppressing_before_completion_stops_delivery() {
let supervisor = SubAgentSupervisor::new(3);
let child = make_session(vec![text_response("child result")]).await;
let agent_id = supervisor
.spawn_with_parent_notification(
child,
"task".to_string(),
"Inspect the module".to_string(),
0,
)
.unwrap();
supervisor.suppress_parent_notification(&agent_id);
supervisor
.wait_with_cancel(&agent_id, &CancellationToken::new())
.await
.unwrap();
assert!(
hub.next_batch(&CancellationToken::new())
supervisor
.next_parent_notification_batch(&CancellationToken::new())
.await
.unwrap()
.is_none()
);
supervisor.shutdown_all().await;
}
#[tokio::test]
async fn closing_a_running_agent_stops_delivery_without_parking_the_parent() {
let supervisor = SubAgentSupervisor::new(3);
let child = make_session(vec![text_response("child result")]).await;
let agent_id = supervisor
.spawn_with_parent_notification(
child,
"task".to_string(),
"Inspect the module".to_string(),
0,
)
.unwrap();
supervisor.close_agent(&agent_id).await.unwrap();
// Must resolve rather than wait for a result that will never arrive.
assert!(
supervisor
.next_parent_notification_batch(&CancellationToken::new())
.await
.unwrap()
.is_none()
);
supervisor.shutdown_all().await;
}
#[tokio::test]