diff --git a/.github/workflows/http2.yaml b/.github/workflows/http2.yaml index 106df919e2..521fb11e48 100644 --- a/.github/workflows/http2.yaml +++ b/.github/workflows/http2.yaml @@ -29,7 +29,7 @@ jobs: strategy: fail-fast: false matrix: - sdk: [dev] + sdk: [stable] steps: - uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd - uses: dart-lang/setup-dart@65eb853c7ba17dde3be364c3d2858773e7144260 diff --git a/pkgs/http2/CHANGELOG.md b/pkgs/http2/CHANGELOG.md index b98deeeefa..f980aea46f 100644 --- a/pkgs/http2/CHANGELOG.md +++ b/pkgs/http2/CHANGELOG.md @@ -1,6 +1,7 @@ ## 3.0.1-wip - Gracefully handle receiving headers on a stream that the client has canceled. (#1799) +- Enforce the locally advertised `SETTINGS_MAX_CONCURRENT_STREAMS` limit on incoming remote streams. ## 3.0.0 diff --git a/pkgs/http2/lib/src/streams/stream_handler.dart b/pkgs/http2/lib/src/streams/stream_handler.dart index 7db1310658..4142ca50c1 100644 --- a/pkgs/http2/lib/src/streams/stream_handler.dart +++ b/pkgs/http2/lib/src/streams/stream_handler.dart @@ -542,6 +542,12 @@ class StreamHandler extends Object with TerminatableMixin, ClosableMixin { ErrorCode.STREAM_CLOSED, ); _closeStreamIdAbnormally(exception.streamId, exception); + } on StreamRefusedException catch (exception) { + _frameWriter.writeRstStreamFrame( + exception.streamId, + ErrorCode.REFUSED_STREAM, + ); + _closeStreamIdAbnormally(exception.streamId, exception); } on StreamException catch (exception) { _frameWriter.writeRstStreamFrame( exception.streamId, @@ -607,6 +613,25 @@ class StreamHandler extends Object with TerminatableMixin, ClosableMixin { if (frame is HeadersFrame) { if (isServer) { + var localLimit = _localSettings.maxConcurrentStreams; + if (localLimit != null) { + // Enforce our own advertised SETTINGS_MAX_CONCURRENT_STREAMS on + // peer-initiated streams. RFC 7540 5.1.2: an endpoint that + // receives a HEADERS frame that causes its advertised concurrent + // stream limit to be exceeded MUST treat this as a stream error + // of type PROTOCOL_ERROR or REFUSED_STREAM. + var activePeerStreams = + _openStreams.values + .where((s) => _isPeerInitiatedStream(s.id)) + .length; + if (activePeerStreams >= localLimit) { + throw StreamRefusedException( + frame.header.streamId, + 'Refusing remote stream: peer exceeded the locally ' + 'advertised SETTINGS_MAX_CONCURRENT_STREAMS ($localLimit).', + ); + } + } var newStream = newRemoteStream(frame.header.streamId); _changeState(newStream, StreamState.Open); diff --git a/pkgs/http2/lib/src/sync_errors.dart b/pkgs/http2/lib/src/sync_errors.dart index 3d11616ad1..6c810609b4 100644 --- a/pkgs/http2/lib/src/sync_errors.dart +++ b/pkgs/http2/lib/src/sync_errors.dart @@ -50,3 +50,11 @@ class StreamClosedException extends StreamException { @override String toString() => 'StreamClosedException(stream id: $streamId): $_message'; } + +class StreamRefusedException extends StreamException { + StreamRefusedException(super.streamId, [super.message = '']); + + @override + String toString() => + 'StreamRefusedException(stream id: $streamId): $_message'; +} diff --git a/pkgs/http2/test/server_test.dart b/pkgs/http2/test/server_test.dart index 2aac5fe2e2..cb7dfb400c 100644 --- a/pkgs/http2/test/server_test.dart +++ b/pkgs/http2/test/server_test.dart @@ -217,6 +217,87 @@ void main() { await Future.wait([serverFun(), clientFun()]); }); }); + + group('max-concurrent-streams', () { + test('exceeding-max-concurrent-streams', () async { + var writeA = StreamController>(); + var writeB = StreamController>(); + + var server = ServerTransportConnection.viaStreams( + writeB.stream, + writeA, + settings: const ServerSettings(concurrentStreamLimit: 2), + ); + + var localSettings = ActiveSettings(); + var clientReader = StreamIterator( + FrameReader(writeA.stream, localSettings).startDecoding(), + ); + + Future nextFrame() async { + expect(await clientReader.moveNext(), true); + return clientReader.current; + } + + var encoder = HPackEncoder(); + var peerSettings = ActiveSettings(); + writeB.add(CONNECTION_PREFACE); + var clientWriter = FrameWriter(encoder, writeB, peerSettings); + + var clientDone = Completer(); + + Future serverFun() async { + var incoming = []; + var subscription = server.incomingStreams.listen(incoming.add); + + await clientDone.future; + + expect(incoming.length, 2); + await subscription.cancel(); + await server.terminate(); + } + + Future clientFun() async { + expect(await nextFrame() is SettingsFrame, true); + clientWriter.writeSettingsAckFrame(); + clientWriter.writeSettingsFrame([]); + expect(await nextFrame() is SettingsFrame, true); + + clientWriter.writeHeadersFrame(1, [ + Header.ascii('a', 'b'), + ], endStream: false); + clientWriter.writeHeadersFrame(3, [ + Header.ascii('a', 'b'), + ], endStream: false); + clientWriter.writeHeadersFrame(5, [ + Header.ascii('a', 'b'), + ], endStream: false); + + var frame = await nextFrame(); + expect( + frame, + isA() + .having( + (f) => f.errorCode, + 'errorCode', + ErrorCode.REFUSED_STREAM, + ) + .having((f) => f.header.streamId, 'header.streamId', 5), + ); + + clientDone.complete(); + + var hasGoaway = await clientReader.moveNext(); + expect(hasGoaway, true); + expect(clientReader.current is GoawayFrame, true); + + var closed = await clientReader.moveNext(); + expect(closed, false); + } + + await [serverFun(), clientFun()].wait; + }); + }); }); }