Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
---
category: Fixed
---
- A `HarnessClient` dropped without an explicit `close()` now closes its notification channel, so a subscription awaiting the next notification returns `Error::TransportClosed` promptly instead of only when the aborted background reader task happens to be torn down. That error reports the closed reason, plus the runtime's exit code and captured stderr tail when they were observed.
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
---
category: Fixed
---
- `Session::run` no longer wedges indefinitely when the runtime dies mid-turn (stdout EOF after `session/prompt` succeeds but before the inbox receipt or root-idle notification arrives). The high-level API surfaces `Error::TransportClosed` with the exit code and stderr tail, as its doc contract always promised ("once the channel (or the client) is closed, `Error::TransportClosed` is returned").
- `Session::run` no longer wedges indefinitely when the runtime dies mid-turn (stdout EOF after `session/prompt` succeeds but before the inbox receipt or root-idle notification arrives). The high-level API surfaces `Error::TransportClosed` with the closed reason and the stderr tail the end-of-stream path captured, plus the exit code when the best-effort poll of the exiting child observed it, as its doc contract always promised ("once the channel (or the client) is closed, `Error::TransportClosed` is returned").
- `NotificationSubscription::recv` parked in `broadcast::Receiver::recv().await` now wakes with `TransportClosed` when the runtime dies spontaneously (stdout EOF without an explicit `close`), instead of hanging forever. The fix shares the original broadcast `Sender` between `HarnessClient` and the read loop and drops it on the read loop's EOF path, so the channel closes and a parked `recv` resolves with `RecvError::Closed`. Subscriptions created after runtime death remain born-failed.
8 changes: 8 additions & 0 deletions src/client/close_ladder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -261,6 +261,14 @@ impl Drop for HarnessClient {
if let Some(handle) = self.stderr_task.take() {
handle.abort();
}
// Take the notification producer as well (the same `Option::take` the
// read loop's EOF path and `finish_teardown` perform): the broadcast
// channel closes here, so a subscription held across a
// drop-without-close wakes with `TransportClosed` deterministically,
// instead of only whenever the aborted read task happens to be
// dropped. Taking it twice is a no-op, so this stays safe when
// `close()` already ran the same tail.
*lock(&self.notifications) = None;
}
}

Expand Down
63 changes: 54 additions & 9 deletions src/client/core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -213,13 +213,14 @@ impl HarnessClient {
let state = Arc::new(Mutex::new(SharedState::default()));
let (notifications, _) = broadcast::channel(broadcast_capacity.max(1));
// The original `Sender` lives inside a shared `NotificationsProducer`
// (an `Arc<Mutex<Option<Sender>>>`) so the read loop's EOF path and
// `close()`'s teardown can both drop it by setting the inner to
// `None`. This closes the broadcast channel: a `recv()` parked in
// `Receiver::recv().await` then resolves with `RecvError::Closed` →
// (an `Arc<Mutex<Option<Sender>>>`) so every teardown path — the read
// loop's EOF tail, `close()`, and `Drop` — can drop it by setting the
// inner to `None`. This closes the broadcast channel: a `recv()` parked
// in `Receiver::recv().await` then resolves with `RecvError::Closed` →
// `Error::TransportClosed` (the documented contract). The read loop
// receives a clone of the `Arc`; subscriptions clone the `Sender` for
// a `Receiver` under the mutex.
// receives a clone of the `Arc`; subscriptions only *subscribe* to the
// shared `Sender` (under the mutex) and keep the resulting `Receiver`,
// never a `Sender` — so no subscription can hold the channel open.
let notifications: NotificationsProducer = Arc::new(Mutex::new(Some(notifications)));

let stdin_shared = Arc::new(tokio::sync::Mutex::new(stdin));
Expand Down Expand Up @@ -506,8 +507,9 @@ impl HarnessClient {
/// [`NotificationSubscription::recv`] rejects immediately with
/// [`Error::TransportClosed`]. The producer is shared with the read loop
/// (see [`NotificationsProducer`](super::NotificationsProducer)), so
/// once the EOF path or `close()` has dropped the `Sender`, the
/// `Option` is `None` and a fresh subscription gets no `Receiver`.
/// once the EOF path, `close()`, or `HarnessClient::drop` has dropped the
/// `Sender`, the `Option` is `None` and a fresh subscription gets no
/// `Receiver`.
pub fn subscribe_session_tree(&self, root: &str) -> NotificationSubscription {
NotificationSubscription {
receiver: lock(&self.notifications)
Expand All @@ -529,7 +531,8 @@ mod tests {
/// Build a `HarnessClient` skeleton carrying only the fields
/// `subscribe_session_tree` consults (the producer, parent map, and
/// state). The task / child / stdin handles are `None`, so dropping the
/// client runs the no-op tail of `Drop` (no background tasks to abort).
/// client aborts no background tasks — but the `Drop` tail still takes the
/// notification producer.
fn client_with_notifications(notifications: NotificationsProducer) -> HarnessClient {
HarnessClient {
child: None,
Expand Down Expand Up @@ -581,4 +584,46 @@ mod tests {
other => panic!("expected TransportClosed, got {other:?}"),
}
}

#[tokio::test]
async fn drop_without_close_takes_the_producer_and_wakes_a_parked_subscription() {
// `Drop` takes the notification producer, so a subscription held
// across a drop-without-close wakes with `TransportClosed`
// deterministically instead of whenever the aborted read task happens
// to be dropped. Without the take the channel would stay open — a
// `Receiver` alone never closes it — and the parked `recv()` would
// hang.
let (tx, _rx) = broadcast::channel::<Notification>(4);
let notifications: NotificationsProducer = Arc::new(Mutex::new(Some(tx)));
let client = client_with_notifications(Arc::clone(&notifications));
let mut subscription = client.subscribe_session_tree("root");

// Park `recv()` on a spawned task so the drop lands while it waits.
let parked = tokio::spawn(async move { subscription.recv().await });
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(
!parked.is_finished(),
"the receiver must still be parked when the client is dropped"
);

drop(client);

assert!(
lock(&notifications).is_none(),
"Drop must take the producer so the broadcast channel closes"
);
match tokio::time::timeout(Duration::from_secs(2), parked).await {
// The variant is the contract; the message text is a constant this
// path sets, so it is not asserted.
Ok(Ok(Err(Error::TransportClosed(_)))) => {}
Ok(Ok(Ok(notification))) => {
panic!("a parked recv delivered a notification: {notification:?}");
}
Ok(Ok(Err(other))) => panic!("expected TransportClosed, got {other:?}"),
Ok(Err(join_err)) => panic!("the parked recv task failed: {join_err}"),
Err(_elapsed) => {
panic!("the parked recv never woke after the client was dropped");
}
}
}
}
22 changes: 13 additions & 9 deletions src/client/subscription.rs
Original file line number Diff line number Diff line change
Expand Up @@ -231,15 +231,19 @@ mod tests {
);
}

/// Regression for the parked-`recv` half of the EOF-death fix. When the
/// last broadcast `Sender` is dropped, the broadcast channel closes and
/// a parked `receiver.recv().await` resolves with `RecvError::Closed`,
/// which `recv()` surfaces as `Error::TransportClosed` with the
/// documented closed reason. The read loop's EOF path now triggers
/// exactly this path by taking the client-owned `Sender` via the
/// shared `NotificationsProducer`; parked-then-dropped is the
/// observation the bug used to miss. If this regresses, the outer
/// 2 s `tokio::time::timeout` fires and the test fails loudly.
/// Pins the receiver-side mechanism of the EOF-death fix: when the last
/// broadcast `Sender` is dropped, the channel closes and a parked
/// `receiver.recv().await` resolves with `RecvError::Closed`, which
/// `recv()` surfaces as `Error::TransportClosed` with the documented
/// closed reason.
///
/// This test owns its channel and drops its own `Sender`, so it exercises
/// the receiver side only: the paths that drop the *client-owned*
/// `Sender` — the read loop's EOF tail, `close()`, and
/// `HarnessClient::drop` — are covered end-to-end in
/// `tests/client_lifecycle.rs::parked_recv_after_runtime_death_returns_closed`
/// and in the `client::core` unit tests. If this regresses, the outer 2 s
/// `tokio::time::timeout` fires and the test fails loudly.
#[tokio::test]
async fn recv_returns_transport_closed_when_channel_closes_while_recv_is_parked() {
let (tx, mut subscription) = subscription_with_capacity(8);
Expand Down
Loading
Loading