Skip to content
Draft
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
4 changes: 4 additions & 0 deletions pkgs/http2/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,10 @@
- `Http2Client` now drops a connection from its pool as soon as it is dead,
instead of keeping it as an idle connection no request may use, and keeps a
healthy connection pooled when the server resets just one of its streams.
- `Http2Client` propagates pausing a response body to the underlying HTTP/2
stream, so a slow reader stops the server at the flow-control window instead
of buffering the whole response in memory, and uses a 4 MiB stream / 16 MiB
connection receive window instead of the protocol's 64 KiB defaults.

## 3.1.0

Expand Down
14 changes: 13 additions & 1 deletion pkgs/http2/lib/src/http2_client.dart
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,13 @@ class Http2Client extends BaseClient {
final transport = ClientConnection(
incoming,
socket,
const ClientSettings(),
// The protocol's default windows (65535 bytes) cap what a peer may have
// in flight at 64 KiB per round trip, per connection. These also bound
// how much a paused response body can hold up (see `_sendOverHttp2`).
const ClientSettings(
streamWindowSize: 4 * 1024 * 1024,
connectionWindowSize: 16 * 1024 * 1024,
),
);
try {
await Future.any([
Expand Down Expand Up @@ -192,7 +198,13 @@ class Http2Client extends BaseClient {

final statusCompleter = Completer<int>();
late final StreamSubscription<StreamMessage> subscription;
// Pausing the body pauses the HTTP/2 stream, which then stops granting
// flow-control credit (WINDOW_UPDATE, RFC 9113 5.2) once what the peer
// was already allowed to send has arrived: a slow reader holds at most
// the stream window in memory rather than the whole response.
final bodyController = StreamController<List<int>>(
onPause: () => subscription.pause(),
onResume: () => subscription.resume(),
onCancel: () {
lease.release();
stream.terminate();
Expand Down
111 changes: 111 additions & 0 deletions pkgs/http2/test/http2_client_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -6,10 +6,15 @@ import 'dart:async';
import 'dart:convert' show ascii;
import 'dart:io';
import 'dart:math';
import 'dart:typed_data';

import 'package:http/http.dart' show ClientException, Request;
import 'package:http2/multiprotocol_server.dart';
import 'package:http2/src/connection_preface.dart';
import 'package:http2/src/frames/frames.dart';
import 'package:http2/src/hpack/hpack.dart';
import 'package:http2/src/http2_client.dart';
import 'package:http2/src/settings/settings.dart';
import 'package:http2/transport.dart';
import 'package:test/test.dart';

Expand Down Expand Up @@ -495,5 +500,111 @@ void main() {
await socket.close();
},
);

test('a-paused-response-body-stops-granting-flow-control-credit', () async {
const streamWindow = 4 * 1024 * 1024;
const connectionWindow = 16 * 1024 * 1024;
const frameSize = 16 * 1024; // The default SETTINGS_MAX_FRAME_SIZE.

// A frame-level server, so the test sees exactly which WINDOW_UPDATE
// frames the client sends, and when.
final context = _serverContext()..setAlpnProtocols(['h2'], true);
final serverSocket = await SecureServerSocket.bind(
'localhost',
0,
context,
);
final accepted = Completer<SecureSocket>();
serverSocket.listen(accepted.complete);

final client = _testClient();
final responseFuture = client.send(
Request('GET', Uri.parse('https://localhost:${serverSocket.port}/')),
);

final socket = await accepted.future;
final writer = FrameWriter(HPackEncoder(), socket, ActiveSettings());
final frames = StreamIterator(
FrameReader(
readConnectionPreface(socket),
ActiveSettings(),
).startDecoding(),
);
Future<Frame> nextFrame() async {
expect(await frames.moveNext(), isTrue, reason: 'connection closed');
return frames.current;
}

// The client only sends its request once it has the server's SETTINGS.
writer.writeSettingsFrame([]);
int? advertisedStreamWindow;
var advertisedConnectionWindow = 65535;
late final int streamId;
while (true) {
final frame = await nextFrame();
if (frame is SettingsFrame && !frame.hasAckFlag) {
for (final setting in frame.settings) {
if (setting.identifier == Setting.SETTINGS_INITIAL_WINDOW_SIZE) {
advertisedStreamWindow = setting.value;
}
}
writer.writeSettingsAckFrame();
} else if (frame is WindowUpdateFrame && frame.header.streamId == 0) {
advertisedConnectionWindow += frame.windowSizeIncrement;
} else if (frame is HeadersFrame) {
streamId = frame.header.streamId;
break;
}
}
expect(advertisedStreamWindow, streamWindow);
expect(advertisedConnectionWindow, connectionWindow);

writer.writeHeadersFrame(streamId, [
Header.ascii(':status', '200'),
], endStream: false);
final response = await responseFuture;
final received = <int>[];
final body = response.stream.listen(received.addAll)..pause();

// Everything the client allows without granting more credit, then a
// PING: the client only answers it once it has processed every frame
// before it.
for (var sent = 0; sent < streamWindow; sent += frameSize) {
writer.writeDataFrame(streamId, Uint8List(frameSize));
}
writer.writePingFrame(1);
while (true) {
final frame = await nextFrame();
if (frame is PingFrame && frame.hasAckFlag) break;
if (frame is WindowUpdateFrame && frame.header.streamId == streamId) {
fail(
'The client granted ${frame.windowSizeIncrement} bytes of credit '
'while the response body was paused.',
);
}
}
expect(received, isEmpty);

// Resuming lets the client consume the window and grant it back.
body.resume();
var granted = 0;
while (granted < streamWindow) {
final frame = await nextFrame();
if (frame is WindowUpdateFrame && frame.header.streamId == streamId) {
granted += frame.windowSizeIncrement;
}
}
expect(granted, streamWindow);

writer.writeDataFrame(streamId, ascii.encode('end'), endStream: true);
await body.asFuture<void>();
expect(received, hasLength(streamWindow + 3));

client.close();
await client.closed;
await frames.cancel();
socket.destroy();
await serverSocket.close();
});
});
}
Loading