From a59c01f8819678e6536c43363d307497f0fe26cf Mon Sep 17 00:00:00 2001 From: Moritz Date: Thu, 8 Oct 2026 10:16:56 +0000 Subject: [PATCH] fix(http2): finish() and terminate() must not hang after a failed socket write MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit > [!NOTE] > This PR was generated by an AI coding agent (Jetski) on behalf of @mosuem. ### Summary `BufferedSink` (the sink behind `FrameWriter`) defined its completion as ```dart _doneFuture = Future.wait([_controller.stream.pipe(dataSink), dataSink.done]); ``` For a `dart:io` `Socket`, `done` only completes after an explicit `close()`. When a write fails — typically `SocketException: Connection reset by peer` because the peer already closed the connection — `addStream` completes with the error, `Stream.pipe` does *not* close the sink, and `Socket.done` never completes. `Future.wait` (non-eager) then waits forever, and so does everything built on `doneFuture`: `FrameWriter.close()`, `Connection.finish()`, `Connection.terminate()`, `ClientPool.terminate()` / `Http2Client.closed`, and grpc-dart's `ClientChannel.shutdown()` (which awaits `transport.finish()`). `pipe` itself already resolves through `dataSink.close()`, i.e. through `done`, in the success case, so the extra wait only ever mattered in the failure case — where it hangs. ### Changes - **`lib/src/async_utils/async_utils.dart`**: `doneFuture` completes once the pipe has finished. A failed write completes it normally (like a cancelled sink already did): the connection is dead and the owner learns about it through the incoming side (`Connection` also terminates itself when `doneFuture` completes). ### Test Verification (Fails Before $\rightarrow$ Passes After) - `buffered-sink-done-after-failed-write` in `test/src/async_utils/async_utils_test.dart`, using a `FailingSink` that behaves like a reset socket (`addStream` fails, `done` never completes without `close()`). - `finish-and-terminate-complete-when-the-socket-write-fails` in `test/client_test.dart`: `ClientTransportConnection.viaStreams(..., FailingSink())`, then `finish()` / `terminate()` must complete. **Before fix:** ```text 00:02 +0 -1: test/src/async_utils/async_utils_test.dart: async_utils buffered-sink-done-after-failed-write [E] TimeoutException after 0:00:02.000000: Future not completed 00:02 +0 -2: test/client_test.dart: client-tests client-errors finish-and-terminate-complete-when-the-socket-write-fails [E] TimeoutException after 0:00:02.000000: Future not completed ``` **After fix:** ```text 00:00 +2: All tests passed! ``` --- pkgs/http2/CHANGELOG.md | 3 ++ .../lib/src/async_utils/async_utils.dart | 17 ++++++-- pkgs/http2/test/client_test.dart | 22 +++++++++++ .../src/async_utils/async_utils_test.dart | 39 +++++++++++++++++++ 4 files changed, 77 insertions(+), 4 deletions(-) diff --git a/pkgs/http2/CHANGELOG.md b/pkgs/http2/CHANGELOG.md index ed5aa1eb30..3cec911a82 100644 --- a/pkgs/http2/CHANGELOG.md +++ b/pkgs/http2/CHANGELOG.md @@ -27,6 +27,9 @@ now extends `StreamTransportException`, as documented for `TransportException` subclasses. `MultiProtocolHttpServer.startServing` forwards an unexpected ALPN protocol to `onError` instead of throwing inside the socket listener. +- `finish()` and `terminate()` no longer hang when a write to the socket fails + (for example because the peer closed the connection); previously they waited + for a `Socket.done` that never completes after a failed write. ## 3.1.0 diff --git a/pkgs/http2/lib/src/async_utils/async_utils.dart b/pkgs/http2/lib/src/async_utils/async_utils.dart index 07aad5e933..03640dfd36 100644 --- a/pkgs/http2/lib/src/async_utils/async_utils.dart +++ b/pkgs/http2/lib/src/async_utils/async_utils.dart @@ -68,10 +68,19 @@ class BufferedSink { // Currently `_doneFuture` will just complete normally if the sink // cancelled. }; - _doneFuture = Future.wait([ - _controller.stream.pipe(dataSink), - dataSink.done, - ]); + // `pipe` completes once `dataSink.close()` has completed, which for + // `dart:io` sinks and `StreamController`s is once `dataSink.done` has. + // + // It must not additionally wait for `dataSink.done`: when a write fails + // (e.g. the peer reset the connection), `addStream` completes with the + // error but `Stream.pipe` does not close the sink, so `dataSink.done` never + // completes - and neither would this future, nor `Connection.finish()` and + // `Connection.terminate()`, which wait for it. + // + // A failed write means the connection is dead; the owner learns about that + // through the incoming side, so like a cancelled sink it just completes + // this future normally. + _doneFuture = _controller.stream.pipe(dataSink).catchError((Object _) {}); } /// The underlying sink. diff --git a/pkgs/http2/test/client_test.dart b/pkgs/http2/test/client_test.dart index 556ae9ea00..4e828e3c35 100644 --- a/pkgs/http2/test/client_test.dart +++ b/pkgs/http2/test/client_test.dart @@ -14,6 +14,7 @@ import 'package:http2/src/settings/settings.dart'; import 'package:http2/transport.dart'; import 'package:test/test.dart'; +import 'src/async_utils/async_utils_test.dart' show FailingSink; import 'src/hpack/hpack_test.dart' show isHeader; void main() { @@ -1693,6 +1694,27 @@ void main() { await Future.wait([serverFun(), clientFun()], eagerError: true); }, ); + + test( + 'finish-and-terminate-complete-when-the-socket-write-fails', + () async { + // Before any of the client's bytes reach the wire, the peer has gone + // away: every write fails and the sink's `done` never completes (as + // with a `dart:io` Socket whose `addStream` failed). + for (final close in [ + (ClientTransportConnection c) => c.finish(), + (ClientTransportConnection c) => c.terminate(), + ]) { + final incoming = StreamController>(); + final connection = ClientTransportConnection.viaStreams( + incoming.stream, + FailingSink(), + ); + await close(connection).timeout(const Duration(seconds: 2)); + await incoming.close(); + } + }, + ); }); }); } diff --git a/pkgs/http2/test/src/async_utils/async_utils_test.dart b/pkgs/http2/test/src/async_utils/async_utils_test.dart index 3917febaa1..98a69353fc 100644 --- a/pkgs/http2/test/src/async_utils/async_utils_test.dart +++ b/pkgs/http2/test/src/async_utils/async_utils_test.dart @@ -3,6 +3,7 @@ // BSD-style license that can be found in the LICENSE file. import 'dart:async'; +import 'dart:io'; import 'package:http2/src/async_utils/async_utils.dart'; import 'package:test/test.dart'; @@ -67,6 +68,15 @@ void main() { ); }); + test('buffered-sink-done-after-failed-write', () async { + final bs = BufferedSink(FailingSink()); + bs.sink.add([1, 2, 3]); + await bs.sink.close(); + // Must not hang on `FailingSink.done`, which (like `Socket.done` after a + // failed `addStream`) never completes. + await bs.doneFuture.timeout(const Duration(seconds: 2)); + }); + test('buffered-bytes-writer', () async { var c = StreamController>(); var writer = BufferedBytesWriter(c); @@ -89,3 +99,32 @@ void main() { }); }); } + +/// Behaves like a `dart:io` [Socket] whose peer has reset the connection: +/// [addStream] fails on the first write and [done] only ever completes through +/// an explicit [close], which `Stream.pipe` does not call after an error. +class FailingSink implements StreamSink> { + final _done = Completer(); + + @override + void add(List data) {} + + @override + void addError(Object error, [StackTrace? stackTrace]) {} + + @override + Future addStream(Stream> stream) async { + await for (final _ in stream) { + throw const SocketException('Connection reset by peer'); + } + } + + @override + Future close() { + if (!_done.isCompleted) _done.complete(); + return _done.future; + } + + @override + Future get done => _done.future; +}