From bd3b27b46c0a74216b39a88aa108c92848998002 Mon Sep 17 00:00:00 2001 From: Alex Li Date: Mon, 20 Jul 2026 20:30:45 +0800 Subject: [PATCH] fix(http2_adapter): handle closed streams during uploads Stop request body subscriptions when the peer closes the HTTP/2 connection, while preserving cancellation deallocation behavior. Co-authored-by: Ivan Elizarov Co-Authored-By: Codex --- plugins/http2_adapter/CHANGELOG.md | 2 + .../http2_adapter/lib/src/http2_adapter.dart | 58 +++++++-- plugins/http2_adapter/test/http2_test.dart | 121 ++++++++++++++++++ 3 files changed, 170 insertions(+), 11 deletions(-) diff --git a/plugins/http2_adapter/CHANGELOG.md b/plugins/http2_adapter/CHANGELOG.md index 4bcb5076e..dbf98f1af 100644 --- a/plugins/http2_adapter/CHANGELOG.md +++ b/plugins/http2_adapter/CHANGELOG.md @@ -5,6 +5,8 @@ See the [Migration Guide][] for the complete breaking changes list.** ## Unreleased +- Prevent uploads from throwing a `StateError` when the HTTP/2 connection closes + before the request body finishes streaming. - Run HTTP integration and certificate pinning tests against the configured test server. diff --git a/plugins/http2_adapter/lib/src/http2_adapter.dart b/plugins/http2_adapter/lib/src/http2_adapter.dart index 3ef5e6644..171eb2582 100644 --- a/plugins/http2_adapter/lib/src/http2_adapter.dart +++ b/plugins/http2_adapter/lib/src/http2_adapter.dart @@ -127,12 +127,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(); @@ -140,16 +134,58 @@ 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? requestSubscription; + final requestCompleter = Completer(); + + void stopRequestStream() { + if (!requestCompleter.isCompleted) { + requestCompleter.complete(); + } + requestSubscription?.cancel().ignore(); + } + + requestSubscription = requestStream!.listen( + (data) { + try { + stream.outgoingMessages.add(DataStreamMessage(data)); + } on StateError { + stopRequestStream(); + } + }, + onError: (Object error, StackTrace stackTrace) { + if (!requestCompleter.isCompleted) { + requestCompleter.completeError(error, stackTrace); + } + }, + onDone: () { + if (!requestCompleter.isCompleted) { + requestCompleter.complete(); + } + }, + cancelOnError: true, + ); + + if (cancelFuture != null) { + final requestSubscriptionWR = WeakReference(requestSubscription); + final requestCompleterWR = WeakReference(requestCompleter); + cancelFuture.whenComplete(() { + final completer = requestCompleterWR.target; + if (completer != null && !completer.isCompleted) { + completer.complete(); + } + requestSubscriptionWR.target?.cancel().ignore(); + streamWR.target?.outgoingMessages.close().ignore(); + }); + } + + Future requestStreamFuture = requestCompleter.future; final sendTimeout = options.sendTimeout ?? Duration.zero; if (sendTimeout > Duration.zero) { requestStreamFuture = requestStreamFuture.timeout( sendTimeout, onTimeout: () { - stream.outgoingMessages.close().catchError((_) {}); + stopRequestStream(); + streamWR.target?.outgoingMessages.close().ignore(); throw DioException.sendTimeout( timeout: sendTimeout, requestOptions: options, diff --git a/plugins/http2_adapter/test/http2_test.dart b/plugins/http2_adapter/test/http2_test.dart index c96f4431e..f77e79831 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() { @@ -108,6 +111,124 @@ void main() { expect(res.data.toString(), contains('TEST')); }); + group('request stream', () { + late ServerSocket serverSocket; + late Http2Adapter adapter; + final serverConnections = []; + + setUp(() async { + serverSocket = await ServerSocket.bind(InternetAddress.loopbackIPv4, 0); + adapter = Http2Adapter(null); + }); + + tearDown(() async { + adapter.close(force: true); + for (final connection in serverConnections) { + await connection.terminate(); + } + serverConnections.clear(); + await serverSocket.close(); + }); + + test( + 'reports a transport error when the connection closes during upload', + () async { + serverSocket.listen((socket) { + final connection = ServerTransportConnection.viaSocket(socket); + serverConnections.add(connection); + connection.incomingStreams.listen((stream) async { + await for (final message in stream.incomingMessages) { + if (message is HeadersStreamMessage) { + await connection.terminate(); + break; + } + } + }); + }); + + final requestStream = (() async* { + for (int i = 0; i < 20; i++) { + await Future.delayed(const Duration(milliseconds: 5)); + yield Uint8List(1024); + } + })(); + final resultCompleter = Completer(); + + void completeResult(Object? result) { + if (!resultCompleter.isCompleted) { + resultCompleter.complete(result); + } + } + + runZonedGuarded( + () async { + try { + await adapter.fetch( + RequestOptions( + path: '/upload', + method: 'POST', + baseUrl: 'http://127.0.0.1:${serverSocket.port}', + ), + requestStream, + null, + ); + completeResult(null); + } catch (error) { + completeResult(error); + } + }, + (error, _) => completeResult(error), + ); + + final result = await resultCompleter.future.timeout( + const Duration(seconds: 2), + ); + expect(result, isA()); + }, + ); + + test('settles the adapter future when an upload is canceled', () async { + final requestDataReceived = Completer(); + serverSocket.listen((socket) { + final connection = ServerTransportConnection.viaSocket(socket); + serverConnections.add(connection); + connection.incomingStreams.listen((stream) async { + await for (final message in stream.incomingMessages) { + if (message is DataStreamMessage && + !requestDataReceived.isCompleted) { + requestDataReceived.complete(); + } + } + stream.sendHeaders( + [Header.ascii(':status', '200')], + endStream: true, + ); + }); + }); + + final requestController = StreamController(); + addTearDown(requestController.close); + final cancelCompleter = Completer(); + final fetchFuture = adapter.fetch( + RequestOptions( + path: '/upload', + method: 'POST', + baseUrl: 'http://127.0.0.1:${serverSocket.port}', + ), + requestController.stream, + cancelCompleter.future, + ); + + requestController.add(Uint8List(1024)); + await requestDataReceived.future.timeout(const Duration(seconds: 2)); + cancelCompleter.complete(); + + final response = await fetchFuture.timeout(const Duration(seconds: 2)); + expect(response.statusCode, 200); + expect(requestController.hasListener, isFalse); + }); + }); + group(ConnectionManager, () { test('returns correct connection', () async { final manager = ConnectionManager();