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
2 changes: 2 additions & 0 deletions plugins/http2_adapter/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
58 changes: 47 additions & 11 deletions plugins/http2_adapter/lib/src/http2_adapter.dart
Original file line number Diff line number Diff line change
Expand Up @@ -127,29 +127,65 @@ class Http2Adapter implements HttpClientAdapter {
final streamWR = WeakReference<ClientTransportStream>(stream);

final hasRequestData = requestStream != null;
if (hasRequestData && cancelFuture != null) {
cancelFuture.whenComplete(() {
streamWR.target?.outgoingMessages.close();
});
}

List<Uint8List>? list;
if (!excludeMethods.contains(options.method) && hasRequestData) {
list = await requestStream.toList();
requestStream = Stream.fromIterable(list);
}

if (hasRequestData) {
Future<dynamic> 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<Uint8List>? requestSubscription;
final requestCompleter = Completer<void>();

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<void> 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,
Expand Down
121 changes: 121 additions & 0 deletions plugins/http2_adapter/test/http2_test.dart
Original file line number Diff line number Diff line change
@@ -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() {
Expand Down Expand Up @@ -108,6 +111,124 @@ void main() {
expect(res.data.toString(), contains('TEST'));
});

group('request stream', () {
late ServerSocket serverSocket;
late Http2Adapter adapter;
final serverConnections = <ServerTransportConnection>[];

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<void>.delayed(const Duration(milliseconds: 5));
yield Uint8List(1024);
}
})();
final resultCompleter = Completer<Object?>();

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<TransportConnectionException>());
},
);

test('settles the adapter future when an upload is canceled', () async {
final requestDataReceived = Completer<void>();
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<Uint8List>();
addTearDown(requestController.close);
final cancelCompleter = Completer<void>();
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();
Expand Down
Loading