Skip to content
Open
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 @@ -19,6 +19,8 @@
- Deliver the complete response instead of a `StreamTransportException` when
the peer sends `RST_STREAM(NO_ERROR)` after `END_STREAM` (RFC 9113
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.

## 3.1.0

Expand Down
14 changes: 14 additions & 0 deletions pkgs/http2/lib/src/flowcontrol/connection_queues.dart
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,20 @@ class ConnectionMessageQueueOut extends Object
});
}

/// Drops the queued [DataMessage]s of the stream [streamId], which is about
/// to be reset.
///
/// Returns `true` if a [HeadersMessage] for [streamId] is still queued, i.e.
/// has not been written to the [FrameWriter] yet. In that case the
/// `RST_STREAM` must be enqueued behind it instead of being written directly:
/// a `RST_STREAM` arriving before the `HEADERS` that open the stream is a
/// connection error for the peer (RFC 9113 section 6.4: "RST_STREAM frames
/// MUST NOT be sent for a stream in the "idle" state").
bool cancelStreamMessages(int streamId) {
_messages.removeWhere((m) => m is DataMessage && m.streamId == streamId);
return _messages.any((m) => m is HeadersMessage && m.streamId == streamId);
}

@override
void onTerminated(Object? error) {
_messages.clear();
Expand Down
17 changes: 16 additions & 1 deletion pkgs/http2/lib/src/streams/stream_handler.dart
Original file line number Diff line number Diff line change
Expand Up @@ -541,7 +541,22 @@ class StreamHandler extends Object with TerminatableMixin, ClosableMixin {
stream.state == StreamState.HalfClosedRemote ||
stream.state == StreamState.ReservedLocal ||
stream.state == StreamState.ReservedRemote) {
_frameWriter.writeRstStreamFrame(stream.id, ErrorCode.CANCEL);
// Drop the stream's queued DATA: the RST_STREAM makes the peer discard
// it anyway. Its queued HEADERS are kept, with the RST_STREAM queued
// behind them instead of written directly: if the HEADERS that open the
// stream have not reached the peer yet (e.g. they are queued behind
// flow-controlled DATA of another stream), a RST_STREAM would arrive for
// a stream the peer considers idle, which is a connection error (RFC
// 9113 sections 5.1 and 6.4). The stream state cannot tell the two
// apart - it advanced when the HEADERS were queued, not written - so all
// queued HEADERS of the stream are treated alike.
if (outgoingQueue.cancelStreamMessages(stream.id)) {
Comment thread
mosuem marked this conversation as resolved.
outgoingQueue.enqueueMessage(
ResetStreamMessage(stream.id, ErrorCode.CANCEL),
);
} else {
_frameWriter.writeRstStreamFrame(stream.id, ErrorCode.CANCEL);
}
_closeStreamAbnormally(stream, null, propagateException: false);
} else if (stream.state == StreamState.Closed &&
!stream.incomingQueue.wasClosed &&
Expand Down
65 changes: 65 additions & 0 deletions pkgs/http2/test/client_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -1446,6 +1446,71 @@ void main() {

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

clientTest('rst-stream-is-sent-after-queued-headers-of-the-stream', (
ClientTransportConnection client,
FrameWriter serverWriter,
StreamIterator<Frame> serverReader,
Future<Frame> Function() nextFrame,
) async {
final handshakeDone = Completer<void>();
final writerBuffers = Completer<void>();
final streamCancelled = Completer<void>();

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

final req1 = await nextFrame() as HeadersFrame;
expect(req1.header.streamId, 1);
// [serverReader] pauses the frame stream right after delivering this
// frame; the pause propagates synchronously back into the client's
// FrameWriter, which now reports that it would buffer. Messages the
// client enqueues from here on stay in its connection queue.
writerBuffers.complete();
await streamCancelled.future;

// Stream 3 was cancelled while its HEADERS were still queued: the
// HEADERS must still reach the wire before the RST_STREAM (RFC 9113
// section 6.4, RST_STREAM on an idle stream is a PROTOCOL_ERROR).
expect(
await nextFrame(),
isA<HeadersFrame>().having((f) => f.header.streamId, 'streamId', 3),
);
expect(
await nextFrame(),
isA<RstStreamFrame>()
.having((f) => f.header.streamId, 'streamId', 3)
.having((f) => f.errorCode, 'errorCode', ErrorCode.CANCEL),
);

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'),
], endStream: true);
await writerBuffers.future;

final s2 = client.makeRequest([Header.ascii(':path', '/s2')]);
s2.terminate();
streamCancelled.complete();

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

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