diff --git a/pkgs/http2/CHANGELOG.md b/pkgs/http2/CHANGELOG.md index f0e3b19085..675ec5bbf7 100644 --- a/pkgs/http2/CHANGELOG.md +++ b/pkgs/http2/CHANGELOG.md @@ -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 diff --git a/pkgs/http2/lib/src/flowcontrol/connection_queues.dart b/pkgs/http2/lib/src/flowcontrol/connection_queues.dart index e315c08bf2..77bb52c01a 100644 --- a/pkgs/http2/lib/src/flowcontrol/connection_queues.dart +++ b/pkgs/http2/lib/src/flowcontrol/connection_queues.dart @@ -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(); diff --git a/pkgs/http2/lib/src/streams/stream_handler.dart b/pkgs/http2/lib/src/streams/stream_handler.dart index 1120e2af03..685ba9f3c5 100644 --- a/pkgs/http2/lib/src/streams/stream_handler.dart +++ b/pkgs/http2/lib/src/streams/stream_handler.dart @@ -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)) { + 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 && diff --git a/pkgs/http2/test/client_test.dart b/pkgs/http2/test/client_test.dart index 344dbb0d2a..19e208d285 100644 --- a/pkgs/http2/test/client_test.dart +++ b/pkgs/http2/test/client_test.dart @@ -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 serverReader, + Future Function() nextFrame, + ) async { + final handshakeDone = Completer(); + final writerBuffers = Completer(); + final streamCancelled = Completer(); + + Future serverFun() async { + expect(await nextFrame(), isA()); + serverWriter.writeSettingsFrame([]); + serverWriter.writeSettingsAckFrame(); + expect(await nextFrame(), isA()); + 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().having((f) => f.header.streamId, 'streamId', 3), + ); + expect( + await nextFrame(), + isA() + .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()); + expect(await serverReader.moveNext(), isFalse); + } + + Future 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(); + await client.finish(); + } + + await Future.wait([serverFun(), clientFun()], eagerError: true); + }); }); }); }