From 7e7a2d271bd7dd24786b7e3423c5382aebd3cce7 Mon Sep 17 00:00:00 2001 From: vanya elizarov Date: Fri, 12 Dec 2025 16:52:36 +0300 Subject: [PATCH 1/3] fix: Http2Adapter does not fail with StateError when connection is closed --- .gitignore | 3 + .../http2_adapter/lib/src/http2_adapter.dart | 55 +++++++++++++---- plugins/http2_adapter/test/http2_test.dart | 60 +++++++++++++++++++ 3 files changed, 106 insertions(+), 12 deletions(-) diff --git a/.gitignore b/.gitignore index 5ca9b8532..de4d65827 100644 --- a/.gitignore +++ b/.gitignore @@ -33,3 +33,6 @@ melos_overrides.yaml # FVM Version Cache .fvm/ .fvmrc + +# Local history +.history \ No newline at end of file diff --git a/plugins/http2_adapter/lib/src/http2_adapter.dart b/plugins/http2_adapter/lib/src/http2_adapter.dart index fd22eb2cf..1710f7f37 100644 --- a/plugins/http2_adapter/lib/src/http2_adapter.dart +++ b/plugins/http2_adapter/lib/src/http2_adapter.dart @@ -126,12 +126,6 @@ class Http2Adapter implements HttpClientAdapter { final streamWR = WeakReference(stream); final hasRequestData = requestStream != null; - if (hasRequestData && cancelFuture != null) { - cancelFuture.whenComplete(() { - streamWR.target?.outgoingMessages.close(); - }); - } - List? list; if (!excludeMethods.contains(options.method) && hasRequestData) { list = await requestStream.toList(); @@ -139,16 +133,49 @@ class Http2Adapter implements HttpClientAdapter { } if (hasRequestData) { - Future requestStreamFuture = requestStream!.listen((data) { - //TODO(EVERYONE): Investigate why this statement can cause "StateError: Bad state: Cannot add event after closing" - stream.outgoingMessages.add(DataStreamMessage(data)); - }).asFuture(); + StreamSubscription? requestSub; + final requestCompleter = Completer(); + + requestSub = requestStream!.listen( + (Uint8List data) { + try { + stream.outgoingMessages.add(DataStreamMessage(data)); + } on StateError { + requestSub?.cancel(); + if (!requestCompleter.isCompleted) { + requestCompleter.complete(); + } + } + }, + onError: (Object e, StackTrace st) { + if (!requestCompleter.isCompleted) { + requestCompleter.completeError(e, st); + } + }, + onDone: () { + if (!requestCompleter.isCompleted) { + requestCompleter.complete(); + } + }, + cancelOnError: true, + ); + + if (cancelFuture != null) { + cancelFuture.whenComplete(() { + requestSub?.cancel().catchError((_) {}).whenComplete(() { + streamWR.target?.outgoingMessages.close().catchError((_) {}); + }); + }); + } + + Future requestStreamFuture = requestCompleter.future; final sendTimeout = options.sendTimeout ?? Duration.zero; if (sendTimeout > Duration.zero) { requestStreamFuture = requestStreamFuture.timeout( sendTimeout, onTimeout: () { - stream.outgoingMessages.close().catchError((_) {}); + requestSub?.cancel().catchError((_) {}); + streamWR.target?.outgoingMessages.close().catchError((_) {}); throw DioException.sendTimeout( timeout: sendTimeout, requestOptions: options, @@ -156,9 +183,13 @@ class Http2Adapter implements HttpClientAdapter { }, ); } + await requestStreamFuture; } - await stream.outgoingMessages.close(); + + try { + await stream.outgoingMessages.close(); + } catch (_) {} final responseSink = StreamController(); final responseHeaders = Headers(); diff --git a/plugins/http2_adapter/test/http2_test.dart b/plugins/http2_adapter/test/http2_test.dart index d8d9d2321..b6e207972 100644 --- a/plugins/http2_adapter/test/http2_test.dart +++ b/plugins/http2_adapter/test/http2_test.dart @@ -1,8 +1,11 @@ import 'dart:async'; +import 'dart:io'; +import 'dart:typed_data'; import 'package:dio/dio.dart'; import 'package:dio_http2_adapter/dio_http2_adapter.dart'; import 'package:dio_test/util.dart'; +import 'package:http2/transport.dart'; import 'package:test/test.dart'; void main() { @@ -99,6 +102,63 @@ void main() { expect(res.data.toString(), contains('TEST')); }); + test( + 'request does not fail with StateError when the server closes the stream before client sends body', + () async { + final serverSocket = + await ServerSocket.bind(InternetAddress.loopbackIPv4, 0); + + serverSocket.listen((rawSocket) { + final serverConn = ServerTransportConnection.viaSocket(rawSocket); + + serverConn.incomingStreams.listen((ServerTransportStream stream) async { + await for (final msg in stream.incomingMessages) { + if (msg is HeadersStreamMessage) { + stream.terminate(); + break; + } + } + }); + }); + + final dio = Dio(); + final adapter = Http2Adapter(null); + dio.httpClientAdapter = adapter; + + final Stream requestStream = (() async* { + for (int i = 0; i < 20; i++) { + await Future.delayed(const Duration(milliseconds: 5)); + yield Uint8List.fromList(List.filled(1024, i)); + } + })(); + + final completer = Completer(); + + runZonedGuarded(() async { + await adapter.fetch( + RequestOptions( + path: '/test', + method: 'POST', + baseUrl: 'http://127.0.0.1:${serverSocket.port}', + headers: {}, + ), + requestStream, + null, + ); + completer.complete(null); + }, (e, _) { + if (!completer.isCompleted) { + completer.complete(e); + } + }); + + final result = await completer.future; + expect(result, isNot(isA())); + + adapter.close(force: true); + await serverSocket.close(); + }); + group(ConnectionManager, () { test('returns correct connection', () async { final manager = ConnectionManager(); From c86ee61df6719440e8051e6bb25f4eb33e660536 Mon Sep 17 00:00:00 2001 From: vanya elizarov Date: Fri, 12 Dec 2025 17:05:28 +0300 Subject: [PATCH 2/3] update changelog --- plugins/http2_adapter/CHANGELOG.md | 1 + 1 file changed, 1 insertion(+) diff --git a/plugins/http2_adapter/CHANGELOG.md b/plugins/http2_adapter/CHANGELOG.md index 8d516f3f5..af4deafa6 100644 --- a/plugins/http2_adapter/CHANGELOG.md +++ b/plugins/http2_adapter/CHANGELOG.md @@ -6,6 +6,7 @@ See the [Migration Guide][] for the complete breaking changes list.** ## Unreleased - Add `handshakeTimeout` (defaults to 15 seconds) to the `ConnectionManager` to prevent long waiting if there's something wrong with the handshake procedure. +- Fix `StateError: Bad state: Cannot add event after closing` caused by race condition e.g. when the server closed the connection before receiving the request body. ## 2.6.0 From 490b2a93cba4bcbc0f94c3fadc7aa025b19576c6 Mon Sep 17 00:00:00 2001 From: vanya elizarov Date: Mon, 15 Dec 2025 15:56:22 +0300 Subject: [PATCH 3/3] fix review comments --- plugins/http2_adapter/CHANGELOG.md | 2 +- plugins/http2_adapter/lib/src/http2_adapter.dart | 6 +++++- plugins/http2_adapter/test/http2_test.dart | 1 + 3 files changed, 7 insertions(+), 2 deletions(-) diff --git a/plugins/http2_adapter/CHANGELOG.md b/plugins/http2_adapter/CHANGELOG.md index af4deafa6..707febef5 100644 --- a/plugins/http2_adapter/CHANGELOG.md +++ b/plugins/http2_adapter/CHANGELOG.md @@ -6,7 +6,7 @@ See the [Migration Guide][] for the complete breaking changes list.** ## Unreleased - Add `handshakeTimeout` (defaults to 15 seconds) to the `ConnectionManager` to prevent long waiting if there's something wrong with the handshake procedure. -- Fix `StateError: Bad state: Cannot add event after closing` caused by race condition e.g. when the server closed the connection before receiving the request body. +- Fix `StateError: Bad state: Cannot add event after closing` caused by race condition e.g. when the server closed the connection before receiving the request body. ## 2.6.0 diff --git a/plugins/http2_adapter/lib/src/http2_adapter.dart b/plugins/http2_adapter/lib/src/http2_adapter.dart index 1710f7f37..c8e193d45 100644 --- a/plugins/http2_adapter/lib/src/http2_adapter.dart +++ b/plugins/http2_adapter/lib/src/http2_adapter.dart @@ -189,7 +189,11 @@ class Http2Adapter implements HttpClientAdapter { try { await stream.outgoingMessages.close(); - } catch (_) {} + } on StateError { + // Ignore StateError, which may occur if the stream is already closed. + } catch (_) { + rethrow; + } final responseSink = StreamController(); final responseHeaders = Headers(); diff --git a/plugins/http2_adapter/test/http2_test.dart b/plugins/http2_adapter/test/http2_test.dart index b6e207972..3d357ff7e 100644 --- a/plugins/http2_adapter/test/http2_test.dart +++ b/plugins/http2_adapter/test/http2_test.dart @@ -115,6 +115,7 @@ void main() { await for (final msg in stream.incomingMessages) { if (msg is HeadersStreamMessage) { stream.terminate(); + serverConn.terminate(); break; } }