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
2 changes: 2 additions & 0 deletions pkgs/http2/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@
Section 8.1).
- Never write a `RST_STREAM` before the still-queued `HEADERS` of the same
stream when a stream is cancelled; drop the stream's queued `DATA` instead.
- Stop head-of-line blocking `HEADERS`, `RST_STREAM`, and `GOAWAY` behind
`DATA` that is waiting for connection-level flow-control credit.

## 3.1.0

Expand Down
83 changes: 65 additions & 18 deletions pkgs/http2/lib/src/flowcontrol/connection_queues.dart
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@ class ConnectionMessageQueueOut extends Object
ensureNotClosingSync(() {
if (!wasTerminated) {
_messages.addLast(message);
if (message is! DataMessage) _nothingSendableWhileBlocked = false;
_trySendMessages();
}
});
Expand All @@ -74,6 +75,7 @@ class ConnectionMessageQueueOut extends Object
/// MUST NOT be sent for a stream in the "idle" state").
bool cancelStreamMessages(int streamId) {
_messages.removeWhere((m) => m is DataMessage && m.streamId == streamId);
_nothingSendableWhileBlocked = false;
return _messages.any((m) => m is HeadersMessage && m.streamId == streamId);
}

Expand All @@ -90,27 +92,73 @@ class ConnectionMessageQueueOut extends Object
}
}

/// Set once a scan of [_messages] under an exhausted connection window has
/// found nothing sendable.
///
/// Enqueueing further [DataMessage]s cannot change that outcome, so the scan
/// (which is linear in the queue length) is skipped until the queue changes
/// in a way that could make a message sendable: a non-DATA message is
/// enqueued, a message is sent, or a stream's messages are cancelled.
/// Without this, every `sendData()` call made while the connection window is
/// exhausted would rescan the whole queue, i.e. quadratic work.
bool _nothingSendableWhileBlocked = false;

/// The next message which can be written now, or `null` if none can.
///
/// Normally that is the first queued message. When the connection-level send
/// window is exhausted and the first message is a [DataMessage], frames which
/// are not subject to flow control (RFC 9113 section 6.9: only DATA frames
/// are) may still be written: `GOAWAY`, and `HEADERS` / `RST_STREAM` of
/// streams that have no earlier queued message. Skipping a message of a
/// stream blocks all later messages of that stream, so the per-stream frame
/// order is preserved.
Message? _peekSendableMessage() {
if (_messages.isEmpty || _frameWriter.bufferIndicator.wouldBuffer) {
return null;
}
final first = _messages.first;
if (first is! DataMessage ||
!_connectionWindow.positiveWindow.wouldBuffer) {
return first;
}
if (_nothingSendableWhileBlocked) return null;

final blockedStreams = <int>{};
for (final message in _messages) {
if (message is GoawayMessage) return message;
if (message is DataMessage || blockedStreams.contains(message.streamId)) {
blockedStreams.add(message.streamId);
} else if (message is PushPromiseMessage) {
// A PUSH_PROMISE must precede any frame of the promised stream and
// the END_STREAM of its associated stream (RFC 9113 section 8.4.1);
// keep it, and everything after it on both streams, in order.
blockedStreams
..add(message.streamId)
..add(message.promisedStreamId);
} else {
return message;
}
}
_nothingSendableWhileBlocked = true;
return null;
}

void _trySendMessages() {
if (!wasTerminated) {
// We can make progress if
// * there is at least one message to send
// * the underlying frame writer / sink / socket doesn't block
// * either one
// * the next message is a non-flow control message (e.g. headers)
// * the next sendable message is a non-flow control message
// * the connection window is positive

if (_messages.isNotEmpty &&
!_frameWriter.bufferIndicator.wouldBuffer &&
(!_connectionWindow.positiveWindow.wouldBuffer ||
_messages.first is! DataMessage)) {
_trySendMessage();
final message = _peekSendableMessage();
if (message != null) {
_trySendMessage(message);

// If we have more messages and we can send them, we'll run them
// using `Timer.run()` to let other things get in-between.
if (_messages.isNotEmpty &&
!_frameWriter.bufferIndicator.wouldBuffer &&
(!_connectionWindow.positiveWindow.wouldBuffer ||
_messages.first is! DataMessage)) {
if (_peekSendableMessage() != null) {
// TODO: If all the frame writer methods would return the
// number of bytes written, we could just say, we loop here until 10kb
// and after words, we'll make `Timer.run()`.
Expand All @@ -122,25 +170,26 @@ class ConnectionMessageQueueOut extends Object
}
}

void _trySendMessage() {
var message = _messages.first;
if (message is HeadersMessage) {
void _trySendMessage(Message message) {
if (identical(message, _messages.first)) {
_messages.removeFirst();
} else {
_messages.remove(message);
}
_nothingSendableWhileBlocked = false;
if (message is HeadersMessage) {
_frameWriter.writeHeadersFrame(
message.streamId,
message.headers,
endStream: message.endStream,
);
} else if (message is PushPromiseMessage) {
_messages.removeFirst();
_frameWriter.writePushPromiseFrame(
message.streamId,
message.promisedStreamId,
message.headers,
);
} else if (message is DataMessage) {
_messages.removeFirst();

if (_connectionWindow.peerWindowSize >= message.bytes.length) {
_connectionWindow.decreaseWindow(message.bytes.length);
_frameWriter.writeDataFrame(
Expand Down Expand Up @@ -170,10 +219,8 @@ class ConnectionMessageQueueOut extends Object
_messages.addFirst(tailMessage);
}
} else if (message is ResetStreamMessage) {
_messages.removeFirst();
_frameWriter.writeRstStreamFrame(message.streamId, message.errorCode);
} else if (message is GoawayMessage) {
_messages.removeFirst();
_frameWriter.writeGoawayFrame(
message.lastStreamId,
message.errorCode,
Expand Down
181 changes: 181 additions & 0 deletions pkgs/http2/test/client_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -1511,6 +1511,187 @@ void main() {

await Future.wait([serverFun(), clientFun()], eagerError: true);
});

clientTest(
'headers-and-rst-not-blocked-or-reordered-on-exhausted-window',
(
ClientTransportConnection client,
FrameWriter serverWriter,
StreamIterator<Frame> serverReader,
Future<Frame> Function() nextFrame,
) async {
final handshakeDone = Completer<void>();
final windowExhausted = Completer<void>();

Future<void> serverFun() async {
serverWriter.writeSettingsFrame([
Setting(Setting.SETTINGS_INITIAL_WINDOW_SIZE, 200000),
]);
expect(await nextFrame(), isA<SettingsFrame>());
serverWriter.writeSettingsAckFrame();
expect(await nextFrame(), isA<SettingsFrame>());
handshakeDone.complete();

final req1 = await nextFrame() as HeadersFrame;
expect(req1.header.streamId, 1);

var s1Bytes = 0;
while (s1Bytes < 65535) {
final data = await nextFrame() as DataFrame;
expect(data.header.streamId, 1);
expect(data.hasEndStreamFlag, isFalse);
s1Bytes += data.bytes.length;
}
expect(s1Bytes, 65535);
windowExhausted.complete();

// Even though the connection send window is exhausted by stream 1,
// stream 3's HEADERS and stream 5's HEADERS + RST_STREAM must not
// be head-of-line blocked, and stream 5's HEADERS must precede its
// RST_STREAM.
expect(
await nextFrame(),
isA<HeadersFrame>()
.having((f) => f.header.streamId, 'streamId', 3)
.having((f) => f.hasEndStreamFlag, 'endStream', isTrue),
);
expect(
await nextFrame(),
isA<HeadersFrame>().having(
(f) => f.header.streamId,
'streamId',
5,
),
);
expect(
await nextFrame(),
isA<RstStreamFrame>()
.having((f) => f.header.streamId, 'streamId', 5)
.having((f) => f.errorCode, 'errorCode', ErrorCode.CANCEL),
);

serverWriter.writeHeadersFrame(3, [
Header.ascii(':status', '200'),
], endStream: true);

// Grant connection window credit so stream 1's tail DATA flushes.
serverWriter.writeWindowUpdate(10000, streamId: 0);
final tail = await nextFrame() as DataFrame;
expect(tail.header.streamId, 1);
expect(tail.bytes.length, 70000 - 65535);
expect(tail.hasEndStreamFlag, isTrue);

serverWriter.writeHeadersFrame(1, [
Header.ascii(':status', '200'),
], endStream: true);

expect(await nextFrame(), isA<GoawayFrame>());
expect(await serverReader.moveNext(), isFalse);
}

Future<void> clientFun() async {
await handshakeDone.future;
final s1 = client.makeRequest([Header.ascii(':path', '/s1')]);
s1.sendData(List<int>.filled(70000, 1), endStream: true);

await windowExhausted.future;

final s2 = client.makeRequest([
Header.ascii(':path', '/s2'),
], endStream: true);
final s3 = client.makeRequest([Header.ascii(':path', '/s3')]);
s3.terminate();

await s2.incomingMessages.drain<void>();
await s1.incomingMessages.drain<void>();
await client.finish();
}

await Future.wait([serverFun(), clientFun()], eagerError: true);
},
);

clientTest(
'trailers-stay-behind-blocked-data-while-other-headers-bypass',
(
ClientTransportConnection client,
FrameWriter serverWriter,
StreamIterator<Frame> serverReader,
Future<Frame> Function() nextFrame,
) async {
final handshakeDone = Completer<void>();
final windowExhausted = Completer<void>();

Future<void> serverFun() async {
serverWriter.writeSettingsFrame([
Setting(Setting.SETTINGS_INITIAL_WINDOW_SIZE, 200000),
]);
expect(await nextFrame(), isA<SettingsFrame>());
serverWriter.writeSettingsAckFrame();
expect(await nextFrame(), isA<SettingsFrame>());
handshakeDone.complete();

final req1 = await nextFrame() as HeadersFrame;
expect(req1.header.streamId, 1);
var s1Bytes = 0;
while (s1Bytes < 65535) {
final data = await nextFrame() as DataFrame;
expect(data.header.streamId, 1);
s1Bytes += data.bytes.length;
}
expect(s1Bytes, 65535);
windowExhausted.complete();

// Stream 1 still has DATA and its trailers queued. Stream 3's
// HEADERS may pass them, stream 1's trailers may not.
expect(
await nextFrame(),
isA<HeadersFrame>().having(
(f) => f.header.streamId,
'streamId',
3,
),
);
serverWriter.writeWindowUpdate(10000, streamId: 0);
final tail = await nextFrame() as DataFrame;
expect(tail.header.streamId, 1);
expect(tail.bytes.length, 70000 - 65535);
expect(tail.hasEndStreamFlag, isFalse);
expect(
await nextFrame(),
isA<HeadersFrame>()
.having((f) => f.header.streamId, 'streamId', 1)
.having((f) => f.hasEndStreamFlag, 'endStream', isTrue),
);

for (final streamId in [1, 3]) {
serverWriter.writeHeadersFrame(streamId, [
Header.ascii(':status', '200'),
], endStream: true);
}
expect(await nextFrame(), isA<GoawayFrame>());
expect(await serverReader.moveNext(), isFalse);
}

Future<void> clientFun() async {
await handshakeDone.future;
final s1 = client.makeRequest([Header.ascii(':path', '/s1')]);
s1.sendData(List<int>.filled(70000, 1));
s1.sendHeaders([Header.ascii('x-trailer', '1')], endStream: true);

await windowExhausted.future;
final s2 = client.makeRequest([
Header.ascii(':path', '/s2'),
], endStream: true);

await s2.incomingMessages.drain<void>();
await s1.incomingMessages.drain<void>();
await client.finish();
}

await Future.wait([serverFun(), clientFun()], eagerError: true);
},
);
});
});
}
Expand Down
Loading