From de200aa2fabb61afa9d3895dfd8a4607b00eb0a3 Mon Sep 17 00:00:00 2001 From: Tang Bohao Date: Mon, 14 Sep 2026 16:49:37 +0800 Subject: [PATCH 1/2] fix(client): close the notification channel when a client is dropped without close() MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `Drop` now takes the notification producer, the same `Option::take` the read loop's EOF tail and `close()`'s `finish_teardown` perform, so a subscription held across a drop-without-close wakes with `Error::TransportClosed` deterministically instead of only whenever the aborted read task happens to be dropped. Test follow-ups on the parked-recv/death suite (from the deep review of #16-#19): - Drop the `exit code: 0` assertions: the read loop attaches the code from a single best-effort poll of the child, so a healthy build may legitimately omit it. Document that instead of asserting it. - Widen the parked scripts to the file's 1 s park convention and the outer bounds to 5 s, so the parked window is reliable and a genuine hang still fails loudly. - Replace the constant-string closed-reason assertions with the behavioural one: the captured stderr tail must reach the parked `recv()` / `Session::run`. - Delete the vacuous `elapsed < 2 s` checks (reaching the `TransportClosed` arm already proves the outer bound was met) and the leftover `eprintln!("CORRECT: …")` verification print. - Reuse `run_prefix("m-1")` for the `Session::run` script instead of an inlined handshake copy. - Add end-to-end coverage for the post-death teardown contracts: a subscription created after a spontaneous EOF death is born-failed, and a second `close()` is a no-op; plus a unit test for the `Drop` wake. Docs: correct the `NotificationsProducer` ownership comment (subscriptions keep a `Receiver`, never a `Sender`) and the subscription test doc (it owns its channel, so the client-side drops are covered elsewhere). --- ...client-drop-closes-notification-channel.md | 5 + .../fix-parked-recv-hang-on-runtime-death.md | 2 +- src/client/close_ladder.rs | 8 + src/client/core.rs | 63 +++++- src/client/subscription.rs | 22 +- tests/client_lifecycle.rs | 190 ++++++++++++------ 6 files changed, 213 insertions(+), 77 deletions(-) create mode 100644 .changes/unreleased/fix-client-drop-closes-notification-channel.md diff --git a/.changes/unreleased/fix-client-drop-closes-notification-channel.md b/.changes/unreleased/fix-client-drop-closes-notification-channel.md new file mode 100644 index 0000000..a5a09cb --- /dev/null +++ b/.changes/unreleased/fix-client-drop-closes-notification-channel.md @@ -0,0 +1,5 @@ +--- +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 (the error reports whatever diagnostics had been observed by then). +- The parked-recv-on-spontaneous-death regression tests no longer report a false regression on a healthy build: the runtime exit code on that path is best-effort (the process may already have been reaped), so they assert only the diagnostics the path always produces — the closed reason and the captured stderr tail — and additionally cover a subscription created after the death and the second `close()` being a no-op. diff --git a/.changes/unreleased/fix-parked-recv-hang-on-runtime-death.md b/.changes/unreleased/fix-parked-recv-hang-on-runtime-death.md index 850e9fe..eab1407 100644 --- a/.changes/unreleased/fix-parked-recv-hang-on-runtime-death.md +++ b/.changes/unreleased/fix-parked-recv-hang-on-runtime-death.md @@ -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. diff --git a/src/client/close_ladder.rs b/src/client/close_ladder.rs index a7efaa3..d4a05eb 100644 --- a/src/client/close_ladder.rs +++ b/src/client/close_ladder.rs @@ -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; } } diff --git a/src/client/core.rs b/src/client/core.rs index bfe5a9f..11dd4c4 100644 --- a/src/client/core.rs +++ b/src/client/core.rs @@ -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>>`) 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>>`) 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)); @@ -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) @@ -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, @@ -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::(4); + let notifications: NotificationsProducer = Arc::new(Mutex::new(Some(tx))); + let client = client_with_notifications(Arc::clone(¬ifications)); + 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(¬ifications).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"); + } + } + } } diff --git a/src/client/subscription.rs b/src/client/subscription.rs index 2b2196b..e076d04 100644 --- a/src/client/subscription.rs +++ b/src/client/subscription.rs @@ -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); diff --git a/tests/client_lifecycle.rs b/tests/client_lifecycle.rs index db9cf56..e8c4a50 100644 --- a/tests/client_lifecycle.rs +++ b/tests/client_lifecycle.rs @@ -20,8 +20,8 @@ use uuid::Uuid; use common::fake_runtime::{ emit, emit_blank, emit_raw, emit_stderr, env_dump_bin, exit, expect, expect_frame, expect_params, fake_runtime_path, fake_runtime_spec, harness_config, ignore_all, respond, - respond_error, server_info_result, sleep_forever_bin, sleep_ms, test_temp_root, test_timeouts, - FakeRuntime, + respond_error, run_prefix, server_info_result, sleep_forever_bin, sleep_ms, test_temp_root, + test_timeouts, FakeRuntime, }; /// The canonical client-side session ids used across scenarios. @@ -766,10 +766,14 @@ impl Drop for ForbiddenEnvGuard { // --- parked-recv / Session::run mid-turn-death regression tests ------------- // -// Three tests guarding the doc contract that `recv()` (and therefore +// Four tests guarding the doc contract that `recv()` (and therefore // `Session::run`) returns `Error::TransportClosed` once the channel or the // client is closed, even when a `recv()` is already parked in -// `broadcast::Receiver::recv().await` at the moment of runtime death. +// `broadcast::Receiver::recv().await` at the moment of runtime death: +// `parked_recv_*`, the `fresh_recv_*` control, `session_run_*`, and +// `post_death_close_*` (the post-death teardown contracts — a second +// `close()` is a no-op, a subscription created after the death is +// born-failed — which rest on the same channel closure). // Regression for the hang introduced in 218b0b3, where the original // broadcast `Sender` was retained on `HarnessClient` for its lifetime, so a // spontaneous death (stdout EOF without `close()`) left the channel open @@ -778,21 +782,26 @@ impl Drop for ForbiddenEnvGuard { /// A single `recv()` future, parked across the death boundary, must return /// `TransportClosed` promptly once the runtime dies mid-turn — not hang. /// -/// The script accepts `session/prompt`, sleeps briefly so the client's -/// `recv()` has time to park in `broadcast::Receiver::recv().await`, then -/// exits 0 (stdout EOF). The read loop's EOF path drops the client-owned -/// `Sender` (mirroring `close()`), the channel closes, and the parked -/// `recv()` resolves with `RecvError::Closed` → `TransportClosed`. -/// Asserted with a 2 s bound so a regression to the hang reports as a -/// `timeout` failure instead of stalling the suite. +/// The script accepts `session/prompt`, sleeps for the file's 1 s park +/// convention so the client's `recv()` is parked in +/// `broadcast::Receiver::recv().await` when the peer writes its last words and +/// exits (stdout EOF). The read loop's EOF path drains stderr and then drops +/// the client-owned `Sender` (mirroring `close()`), the channel closes, and +/// the parked `recv()` resolves with `RecvError::Closed` → `TransportClosed` +/// carrying the captured stderr tail. The 5 s outer bound keeps a regression +/// to the hang loud (a `timeout` failure) with ample headroom over the park. #[tokio::test] async fn parked_recv_after_runtime_death_returns_closed() { let mut rt = FakeRuntime::spawn(&[ expect("session/prompt"), respond(json!({"messageId": "m-1"})), - // Long enough for the client to park in recv(); short enough to - // bound the test (the deadline below is 2 s). - sleep_ms(300), + // Long enough for the client to park in recv() (the 1 s park + // convention used above); short enough to bound the test (the + // deadline below is 5 s). + sleep_ms(1000), + // The runtime's last words: the stderr task captures them and the + // EOF path embeds the tail in the closed error. + emit_stderr("fatal: runtime panicked"), // stdout EOF closes here; the read loop's EOF path drops the // client-owned broadcast Sender. exit(0), @@ -808,25 +817,22 @@ async fn parked_recv_after_runtime_death_returns_closed() { assert_eq!(message_id, "m-1"); let started = std::time::Instant::now(); - let outcome = tokio::time::timeout(Duration::from_secs(2), subscription.recv()).await; + let outcome = tokio::time::timeout(Duration::from_secs(5), subscription.recv()).await; let elapsed = started.elapsed(); match outcome { Ok(Err(Error::TransportClosed(message))) => { + // The stderr tail is guaranteed on this path only because the + // read loop drains stderr (bounded) *before* taking the producer + // at `src/client/read_loop.rs:86-107` — the ordering that lets + // the parked `recv()` wake with the diagnostics. The exit code is + // deliberately NOT asserted: it comes from a single best-effort + // poll of the child (`src/client/read_loop.rs:71-81`) that may + // find the child not yet reaped or the lock held, so a healthy + // build legitimately omits `exit code:` from this message. assert!( - elapsed < Duration::from_secs(2), - "recv should resolve on its own once the runtime dies, not via the outer timeout" - ); - assert!( - message.contains("DeepSeek Harness runtime closed"), - "the closed reason rides in the diagnostics: {message}" - ); - // The exit code is the best-effort poll the read loop performs - // on EOF (the peer exited 0); seeing it proves the EOF-path - // diagnostics land in the parked-recv error. - assert!( - message.contains("exit code: 0"), - "the EOF-path exit code must reach the parked recv: {message}" + message.contains("fatal: runtime panicked"), + "the EOF-path stderr tail must reach the parked recv: {message}" ); } Ok(Ok(notification)) => { @@ -837,7 +843,7 @@ async fn parked_recv_after_runtime_death_returns_closed() { } Err(_elapsed) => { panic!( - "recv hung for 2 s after runtime death — the parked-recv-on-death bug is back \ + "recv hung for 5 s after runtime death — the parked-recv-on-death bug is back \ (elapsed = {elapsed:?}); see `NotificationSubscription::recv` and the \ read-loop EOF path" ); @@ -877,10 +883,7 @@ async fn fresh_recv_after_runtime_death_returns_closed() { let deadline = std::time::Instant::now() + Duration::from_secs(2); loop { match tokio::time::timeout(Duration::from_millis(50), subscription.recv()).await { - Ok(Err(Error::TransportClosed(message))) => { - eprintln!("CORRECT: fresh recv after death returned TransportClosed: {message}"); - break; - } + Ok(Err(Error::TransportClosed(_))) => break, Ok(Ok(notification)) => { panic!( "fresh recv returned a notification after the runtime had died: \ @@ -911,24 +914,28 @@ async fn fresh_recv_after_runtime_death_returns_closed() { /// before the inbox receipt / root idle arrives). The high-level API's /// Phase 1 / Phase 2 waits park in `subscription.recv().await` (documented /// unbounded); the fix lets the parked `recv()` resolve via the channel -/// closing so `Session::run` propagates `TransportClosed`. +/// closing so `Session::run` propagates `TransportClosed` with the captured +/// stderr tail. The exit code is deliberately NOT asserted — it comes from +/// the same best-effort child poll the parked-recv test documents, so a +/// healthy build may legitimately omit it. #[tokio::test] async fn session_run_surfaces_transport_closed_on_mid_turn_death() { - let config = harness_config(&[ - expect_params( - "initialize", - json!({"provider": "deepseek-official", "model": "deepseek-v4-flash"}), - ), - respond(server_info_result()), - expect("session/prompt"), - respond(json!({"messageId": "m-1"})), + // The parity handshake comes from the shared helper (the pinned + // provider/model defaults live there and nowhere else). + let mut script = run_prefix("m-1"); + script.extend([ // The peer dies *after* `session/prompt` acknowledges and *before* - // any notification is emitted — Phase 1 has parked in recv(). - sleep_ms(300), + // any notification is emitted — Phase 1 has parked in recv() (the + // file's 1 s park convention); short enough to bound the test (the + // deadline below is 5 s). + sleep_ms(1000), + // The runtime's last words: the stderr task captures them and the + // EOF path embeds the tail in the closed error. + emit_stderr("fatal: runtime panicked"), // stdout EOF closes here. exit(0), - ]) - .expect("serialize script"); + ]); + let config = harness_config(&script).expect("serialize script"); // Bound the `Session::run` call so a regression to the hang fails // loudly as a timeout instead of stalling the suite. (Production code @@ -941,7 +948,7 @@ async fn session_run_surfaces_transport_closed_on_mid_turn_death() { let started = std::time::Instant::now(); let outcome = tokio::time::timeout( - Duration::from_secs(2), + Duration::from_secs(5), session.run(Input::Text("hello".to_string()), None), ) .await; @@ -949,17 +956,12 @@ async fn session_run_surfaces_transport_closed_on_mid_turn_death() { match outcome { Ok(Err(Error::TransportClosed(message))) => { + // Same guarantee as the parked-recv test: the read loop drains + // stderr (bounded) *before* taking the producer, so the tail the + // EOF path captured reaches the error `Session::run` surfaces. assert!( - elapsed < Duration::from_secs(2), - "Session::run should surface the death on its own, not via the outer timeout" - ); - assert!( - message.contains("DeepSeek Harness runtime closed"), - "the closed reason must ride in the diagnostics: {message}" - ); - assert!( - message.contains("exit code: 0"), - "the EOF-path exit code must reach Session::run: {message}" + message.contains("fatal: runtime panicked"), + "the EOF-path stderr tail must reach Session::run: {message}" ); } Ok(Ok(result)) => { @@ -970,7 +972,7 @@ async fn session_run_surfaces_transport_closed_on_mid_turn_death() { } Err(_elapsed) => { panic!( - "Session::run hung for 2 s after runtime death — the mid-turn-death wedge \ + "Session::run hung for 5 s after runtime death — the mid-turn-death wedge \ is back (elapsed = {elapsed:?}); see the parked recv in Phase 1/2 and the \ read-loop EOF path" ); @@ -979,3 +981,75 @@ async fn session_run_surfaces_transport_closed_on_mid_turn_death() { harness.close().await.expect("close reaps the dead peer"); } + +/// Post-death teardown through the public API, with no `close()` in between: +/// after a spontaneous stdout EOF the read loop owns the teardown, so a +/// subscription created *after* the death is born-failed (`recv()` returns +/// `TransportClosed` immediately instead of parking), and `close()` still +/// reaps cleanly — twice, since a second `close()` is a documented no-op +/// returning `Ok(())`. +/// +/// This is the end-to-end guard for contracts the hand-built client +/// skeletons in `src/client/core.rs` can only approximate: the producer take +/// at `src/client/read_loop.rs:107` is what makes both hold. +#[tokio::test] +async fn post_death_close_is_a_noop_and_subscriptions_are_born_failed() { + let mut rt = FakeRuntime::spawn(&[ + expect("session/prompt"), + respond(json!({"messageId": "m-1"})), + // Long enough for the client to park in recv() (the file's 1 s park + // convention). + sleep_ms(1000), + // stdout EOF closes here. + exit(0), + ]) + .expect("spawn fake runtime"); + + let message_id = rt + .client + .session_prompt(ROOT_SESSION, vec![ContentBlock::Text { text: "hi".into() }]) + .await + .expect("session/prompt succeeds"); + assert_eq!(message_id, "m-1"); + + // Wait for the read loop's EOF tail before asserting on its effects: a + // fresh `recv()` re-checks `state.closed` at the top of every call, so + // polling it until it reports the death is what makes the post-death + // assertions deterministic (the peer is dead either way — this just stops + // racing the tail). + let mut subscription = rt.client.subscribe_session_tree(ROOT_SESSION); + let deadline = std::time::Instant::now() + Duration::from_secs(5); + loop { + match tokio::time::timeout(Duration::from_millis(50), subscription.recv()).await { + Ok(Err(Error::TransportClosed(_))) => break, + Ok(Ok(notification)) => { + panic!("a notification arrived after the runtime had died: {notification:?}"); + } + Ok(Err(other)) => panic!("unexpected error after the runtime died: {other:?}"), + Err(_elapsed) => assert!( + std::time::Instant::now() < deadline, + "the read loop never reported the runtime death on EOF" + ), + } + } + + // Born-failed: the EOF tail took the producer, so a subscription created + // now carries no receiver and must reject instead of parking. + let mut after_death = rt.client.subscribe_session_tree(ROOT_SESSION); + match tokio::time::timeout(Duration::from_secs(5), after_death.recv()).await { + // `TransportClosed` is the whole contract here; the message text is a + // constant the client sets on this path, so it is not asserted. + Ok(Err(Error::TransportClosed(_))) => {} + Ok(Ok(notification)) => { + panic!("a born-failed subscription delivered a notification: {notification:?}"); + } + Ok(Err(other)) => panic!("a born-failed subscription must reject, got {other:?}"), + Err(_elapsed) => panic!("a born-failed subscription must not park in recv()"), + } + + // Death without `close()` leaves the client in a consistent terminal + // state: the first `close()` reaps the already-dead peer, and a second + // one is a no-op. + rt.client.close().await.expect("close reaps the dead peer"); + rt.client.close().await.expect("a second close is a no-op"); +} From 735ee5fedb9818d1be6b67ef3b038ceda671f511 Mon Sep 17 00:00:00 2001 From: Tang Bohao Date: Mon, 14 Sep 2026 17:36:11 +0800 Subject: [PATCH 2/2] test(client): restore the closed-reason assertion and de-rot comment citations The closed-reason substring is not best-effort: every `TransportClosed` branch of `NotificationSubscription::recv` builds its error from the same constant, and `Session::run` propagates the error unchanged. The assertion removed alongside the unsound exit-code one was therefore sound, and is restored in the parked-path arms of both end-to-end tests. The exit-code assertions stay removed (the read loop attaches the code from a single best-effort poll of the child, whose emission is conditional), as do the `elapsed < 2s` checks, which only restated the outer timeout that gates the arm. Hard-coded `read_loop.rs` line numbers in the test comments are replaced by symbol references ("the bounded stderr drain", "the best-effort exit poll", "the producer take") so the comments cannot rot. The drop-without-close changelog fragment no longer describes the SDK's own test methodology: the consumer-relevant half (the diagnostics carried by the error) folds into the remaining bullet. --- ...client-drop-closes-notification-channel.md | 3 +- tests/client_lifecycle.rs | 37 ++++++++++++------- 2 files changed, 25 insertions(+), 15 deletions(-) diff --git a/.changes/unreleased/fix-client-drop-closes-notification-channel.md b/.changes/unreleased/fix-client-drop-closes-notification-channel.md index a5a09cb..54b5cc3 100644 --- a/.changes/unreleased/fix-client-drop-closes-notification-channel.md +++ b/.changes/unreleased/fix-client-drop-closes-notification-channel.md @@ -1,5 +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 (the error reports whatever diagnostics had been observed by then). -- The parked-recv-on-spontaneous-death regression tests no longer report a false regression on a healthy build: the runtime exit code on that path is best-effort (the process may already have been reaped), so they assert only the diagnostics the path always produces — the closed reason and the captured stderr tail — and additionally cover a subscription created after the death and the second `close()` being a no-op. +- 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. diff --git a/tests/client_lifecycle.rs b/tests/client_lifecycle.rs index e8c4a50..fdd127e 100644 --- a/tests/client_lifecycle.rs +++ b/tests/client_lifecycle.rs @@ -822,14 +822,19 @@ async fn parked_recv_after_runtime_death_returns_closed() { match outcome { Ok(Err(Error::TransportClosed(message))) => { - // The stderr tail is guaranteed on this path only because the - // read loop drains stderr (bounded) *before* taking the producer - // at `src/client/read_loop.rs:86-107` — the ordering that lets - // the parked `recv()` wake with the diagnostics. The exit code is - // deliberately NOT asserted: it comes from a single best-effort - // poll of the child (`src/client/read_loop.rs:71-81`) that may - // find the child not yet reaped or the lock held, so a healthy - // build legitimately omits `exit code:` from this message. + // Every closed path builds its error from the same constant + // reason, so the reason is guaranteed here. The stderr tail is + // guaranteed too, because the read loop drains stderr (bounded) + // *before* the producer take — the ordering that lets the parked + // `recv()` wake with the diagnostics. The exit code is + // deliberately NOT asserted: it comes from the read loop's single + // best-effort poll of the child, which may find it not yet reaped + // or the lock held, so a healthy build legitimately omits + // `exit code:` from this message. + assert!( + message.contains("DeepSeek Harness runtime closed"), + "the closed reason must reach the parked recv: {message}" + ); assert!( message.contains("fatal: runtime panicked"), "the EOF-path stderr tail must reach the parked recv: {message}" @@ -956,9 +961,15 @@ async fn session_run_surfaces_transport_closed_on_mid_turn_death() { match outcome { Ok(Err(Error::TransportClosed(message))) => { - // Same guarantee as the parked-recv test: the read loop drains - // stderr (bounded) *before* taking the producer, so the tail the - // EOF path captured reaches the error `Session::run` surfaces. + // Same guarantees as the parked-recv test: the reason riding + // every closed error is the shared constant, and the read loop + // drains stderr (bounded) *before* the producer take, so the tail + // the EOF path captured reaches the error `Session::run` + // surfaces. + assert!( + message.contains("DeepSeek Harness runtime closed"), + "the closed reason must reach Session::run: {message}" + ); assert!( message.contains("fatal: runtime panicked"), "the EOF-path stderr tail must reach Session::run: {message}" @@ -990,8 +1001,8 @@ async fn session_run_surfaces_transport_closed_on_mid_turn_death() { /// returning `Ok(())`. /// /// This is the end-to-end guard for contracts the hand-built client -/// skeletons in `src/client/core.rs` can only approximate: the producer take -/// at `src/client/read_loop.rs:107` is what makes both hold. +/// skeletons in `src/client/core.rs` can only approximate: the read loop's +/// producer take on the EOF path is what makes both hold. #[tokio::test] async fn post_death_close_is_a_noop_and_subscriptions_are_born_failed() { let mut rt = FakeRuntime::spawn(&[