diff --git a/pkgs/http2/CHANGELOG.md b/pkgs/http2/CHANGELOG.md index b98deeeefa..ddb45d513c 100644 --- a/pkgs/http2/CHANGELOG.md +++ b/pkgs/http2/CHANGELOG.md @@ -1,6 +1,8 @@ ## 3.0.1-wip - Gracefully handle receiving headers on a stream that the client has canceled. (#1799) +- Add `Http2Client` (`package:http2/client.dart`), a pooled, multiplexed + `package:http` `Client` backed by HTTP/2 connections. ## 3.0.0 diff --git a/pkgs/http2/README.md b/pkgs/http2/README.md index edb75c99d6..fce690f59e 100644 --- a/pkgs/http2/README.md +++ b/pkgs/http2/README.md @@ -53,3 +53,27 @@ Future main() async { An example with better error handling is available [here][example]. See the [API docs][api] for more details. + +## Pooled `http.Client` + +`package:http2/client.dart` provides `Http2Client`, a `package:http` +`Client` that pools and multiplexes requests over shared HTTP/2 connections +instead of opening one connection per request. This is useful for workloads +that send many concurrent requests to the same host or hosts, where +`dart:io`'s `HttpClient` (HTTP/1.1 only) would otherwise open a new TCP+TLS +connection per request. + +```dart +import 'package:http2/client.dart'; + +Future main() async { + final client = Http2Client(); + final response = await client.get(Uri.parse('https://example.com/')); + print(response.body); + await client.terminate(); +} +``` + +A connection is dialed per `host:port` as needed, so a single `Http2Client` +is safe to reuse across requests to different hosts. See the example +[here](example/pooled_client.dart). diff --git a/pkgs/http2/example/pooled_client.dart b/pkgs/http2/example/pooled_client.dart new file mode 100644 index 0000000000..97b1951a0e --- /dev/null +++ b/pkgs/http2/example/pooled_client.dart @@ -0,0 +1,34 @@ +// Copyright (c) 2026, the Dart project authors. Please see the AUTHORS file +// for details. All rights reserved. Use of this source code is governed by a +// BSD-style license that can be found in the LICENSE file. + +import 'dart:io'; + +import 'package:http2/client.dart'; + +/// Sends several concurrent requests through a single [Http2Client], +/// demonstrating that they share pooled HTTP/2 connections instead of each +/// opening their own. +void main(List args) async { + if (args.length != 1) { + print('Usage: dart pooled_client.dart '); + exit(1); + } + + final uri = Uri.parse(args[0]); + final client = Http2Client(); + + try { + final responses = await Future.wait( + List.generate(5, (_) => client.get(uri)), + ); + for (final response in responses) { + print('${response.statusCode}: ${response.body.length} bytes'); + } + print('Connections used: ${client.connectionCount}'); + } finally { + // Waits for the requests above to finish before closing every + // connection - see Http2Client.terminate(). + await client.terminate(); + } +} diff --git a/pkgs/http2/lib/client.dart b/pkgs/http2/lib/client.dart new file mode 100644 index 0000000000..90759b1994 --- /dev/null +++ b/pkgs/http2/lib/client.dart @@ -0,0 +1,13 @@ +// Copyright (c) 2026, the Dart project authors. Please see the AUTHORS file +// for details. All rights reserved. Use of this source code is governed by a +// BSD-style license that can be found in the LICENSE file. + +/// A pooled, multiplexed `package:http` `Client` backed by HTTP/2 +/// connections. +/// +/// See [Http2Client]. +library; + +import 'src/http2_client.dart' show Http2Client; + +export 'src/http2_client.dart' show Http2Client; diff --git a/pkgs/http2/lib/src/client_pool.dart b/pkgs/http2/lib/src/client_pool.dart new file mode 100644 index 0000000000..b33e99c3f9 --- /dev/null +++ b/pkgs/http2/lib/src/client_pool.dart @@ -0,0 +1,240 @@ +// Copyright (c) 2026, the Dart project authors. Please see the AUTHORS file +// for details. All rights reserved. Use of this source code is governed by a +// BSD-style license that can be found in the LICENSE file. + +import 'dart:async'; +import 'dart:math'; + +import 'package:collection/collection.dart'; +import 'package:meta/meta.dart'; + +class _PooledResource { + _PooledResource(this.future); + final Future future; + int inFlight = 0; + bool failed = false; + + /// The resolved value of [future], once available - kept so scheduling can + /// consult the resource itself from paths that are synchronous by design. + T? value; +} + +/// A claim on one concurrency slot of a pooled resource. +/// +/// Held from [ClientPool.acquire] until [release], which lets a slot outlive +/// the future that produced it - an HTTP/2 response, for instance, is returned +/// as soon as its headers arrive but keeps its stream open until the body ends. +class PoolLease { + PoolLease._(this._pool, this._resource, this.value); + + final ClientPool _pool; + final _PooledResource _resource; + + /// The resource this slot was claimed on. + final T value; + + var _released = false; + + /// Stops the pool routing new work to this resource. + void markFailed() => _resource.failed = true; + + /// Gives the slot back. Idempotent, so it is safe to call from several + /// terminal paths that may race. + void release() { + if (_released) return; + _released = true; + _pool._release(_resource); + } +} + +/// A pool of resources of type [T]. +/// +/// Packs load onto the most-full resource under its capacity (rather than +/// spreading evenly across resources), opens a new resource once existing +/// ones are full, stops routing new work to a resource once an operation on +/// it throws, and garbage-collects idle resources past [maxIdleResources]. +/// +/// Capacity is [maxConcurrentOperations], lowered to whatever limit a +/// resource reports for itself via `concurrencyLimitOf`. +class ClientPool { + ClientPool( + Future Function() create, { + required this.maxConcurrentOperations, + required Future Function(T resource) destroy, + this.maxIdleResources = 1, + int? Function(T resource)? concurrencyLimitOf, + }) : _create = create, + _destroy = destroy, + _concurrencyLimitOf = concurrencyLimitOf; + + final Future Function() _create; + final Future Function(T resource) _destroy; + + /// Reports a resource's own concurrency limit, or `null` if it imposes none. + /// + /// Consulted on every scheduling decision rather than cached, so a limit + /// the resource revises over its lifetime is picked up. + final int? Function(T resource)? _concurrencyLimitOf; + + final int maxConcurrentOperations; + final int maxIdleResources; + + final _resources = <_PooledResource>[]; + final _pendingDestroys = >{}; + var _terminated = false; + Completer? _drained; + Future? _termination; + + /// The number of resources currently in the pool. + int get size => _resources.length; + + /// The number of in-flight operations across every resource. For testing. + @visibleForTesting + int get opCount => _resources.map((resource) => resource.inFlight).sum; + + /// Claims a slot on an available (or newly created) resource. + /// + /// The caller owns the returned lease and must [PoolLease.release] it on + /// every path, errors included, or the slot leaks and [terminate] never + /// drains. Prefer [run] unless the slot has to outlive the future that + /// produced whatever the caller is returning. + Future> acquire() async { + if (_terminated) { + throw StateError('This pool has already been terminated.'); + } + + while (true) { + final pooled = _acquire(); + pooled.inFlight++; + final T value; + try { + value = pooled.value ??= await pooled.future; + } catch (_) { + pooled.failed = true; + _release(pooled); + rethrow; + } + + if (pooled.inFlight <= _capacityOf(pooled)) { + return PoolLease._(this, pooled, value); + } + _release(pooled); + } + } + + void _release(_PooledResource resource) { + resource.inFlight--; + if (_terminated) { + _maybeCompleteDrain(); + } else { + _collectIfIdle(resource); + } + } + + /// Runs [operation] on an available (or newly created) resource. + Future run(Future Function(T resource) operation) async { + final lease = await acquire(); + try { + return await operation(lease.value); + } catch (_) { + lease.markFailed(); + rethrow; + } finally { + lease.release(); + } + } + + // Synchronous (no `await`), so concurrent calls can't race each other + // into both creating a resource before either sees the other's. + _PooledResource _acquire() { + _PooledResource? selected; + for (final resource in _resources) { + if (resource.failed) continue; + if (resource.inFlight < _capacityOf(resource) && + (selected == null || resource.inFlight > selected.inFlight)) { + selected = resource; + } + } + if (selected != null) return selected; + + final resource = _PooledResource(_create()); + _resources.add(resource); + return resource; + } + + /// How many operations [resource] may run at once. + /// + /// A resource that hasn't been created yet, or that reports no limit of its + /// own, is held to [maxConcurrentOperations]. + int _capacityOf(_PooledResource resource) { + final value = resource.value; + final limitOf = _concurrencyLimitOf; + if (value == null || limitOf == null) return maxConcurrentOperations; + final limit = limitOf(value); + return limit == null + ? maxConcurrentOperations + : max(1, min(maxConcurrentOperations, limit)); + } + + void _collectIfIdle(_PooledResource resource) { + if (resource.inFlight > 0) return; + if (!resource.failed && !_hasExcessIdleCapacity(resource)) return; + + _resources.remove(resource); + _startDestroy(resource); + } + + /// Starts destroying [resource] without waiting for it, so that an operation + /// is never held up by its resource's teardown - `destroy` may itself wait + /// on unrelated work still running on that resource. + /// + /// Tracked in [_pendingDestroys] so [terminate] can still promise that every + /// resource has actually been destroyed by the time it completes. + void _startDestroy(_PooledResource resource) { + final value = resource.value; + final done = + value != null ? _destroy(value) : resource.future.then(_destroy); + final tracked = done.catchError((Object _) {}); + _pendingDestroys.add(tracked); + unawaited(tracked.whenComplete(() => _pendingDestroys.remove(tracked))); + } + + bool _hasExcessIdleCapacity(_PooledResource resource) { + final idleCapacity = + _resources + .map( + (other) => other.failed ? 0 : _capacityOf(other) - other.inFlight, + ) + .sum; + return idleCapacity > maxIdleResources * _capacityOf(resource); + } + + void _maybeCompleteDrain() { + if (_drained case final drained? + when !drained.isCompleted && opCount == 0) { + drained.complete(); + } + } + + /// Waits for in-flight operations to finish, then destroys every + /// resource in the pool. No further operations can run afterward. + /// + /// Idempotent: concurrent and repeated calls all observe the same shutdown. + Future terminate() => _termination ??= _terminate(); + + Future _terminate() async { + _terminated = true; + + if (opCount > 0) { + _drained = Completer(); + await _drained!.future; + } + + final resources = _resources.toList(); + _resources.clear(); + for (final resource in resources) { + _startDestroy(resource); + } + await Future.wait(_pendingDestroys.toList()); + } +} diff --git a/pkgs/http2/lib/src/connection.dart b/pkgs/http2/lib/src/connection.dart index 6a9f3fc9de..3a5e12abe0 100644 --- a/pkgs/http2/lib/src/connection.dart +++ b/pkgs/http2/lib/src/connection.dart @@ -513,6 +513,28 @@ class ClientConnection extends Connection implements ClientTransportConnection { bool get isOpen => !_state.isFinishing && !_state.isTerminated && _streams.canOpenStream; + /// Whether this connection is shutting down or already dead. + /// + /// This is the half of [isOpen] that says the connection is unusable for + /// good, as opposed to [canOpenStream]'s "not right now". Separating them + /// lets a caller retry on a busy connection without discarding it. + bool get isClosing => _state.isFinishing || _state.isTerminated; + + /// Whether another concurrent stream may be opened right now, per the limit + /// the peer has advertised. Says nothing about whether the connection is + /// still alive - see [isClosing]. + bool get canOpenStream => _streams.canOpenStream; + + /// The maximum number of concurrent streams the peer currently allows, per + /// its most recent SETTINGS_MAX_CONCURRENT_STREAMS (RFC 7540 6.5.2), or + /// `null` if it hasn't advertised a limit. + /// + /// Deliberately not on [ClientTransportConnection]: that class can only be + /// subtyped with `implements`, so adding a member to it would break every + /// downstream implementation. + int? get peerMaxConcurrentStreams => + _settingsHandler.peerSettings.maxConcurrentStreams; + @override ClientTransportStream makeRequest( List
headers, { diff --git a/pkgs/http2/lib/src/http2_client.dart b/pkgs/http2/lib/src/http2_client.dart new file mode 100644 index 0000000000..c1d1743667 --- /dev/null +++ b/pkgs/http2/lib/src/http2_client.dart @@ -0,0 +1,444 @@ +// Copyright (c) 2026, the Dart project authors. Please see the AUTHORS file +// for details. All rights reserved. Use of this source code is governed by a +// BSD-style license that can be found in the LICENSE file. + +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; +import 'dart:typed_data'; + +import 'package:collection/collection.dart'; +import 'package:http/http.dart'; +import 'package:pool/pool.dart'; + +import '../transport.dart'; +import 'client_pool.dart'; +import 'connection.dart'; + +/// A pooled, multiplexed `http.Client` backed by HTTP/2 connections. +/// +/// Every request is sent as its own HTTP/2 stream on a shared connection - +/// see [ClientPool] - dialed per `host:port` via +/// `SecureSocket.connect(..., supportedProtocols: ['h2'])`. Once a connection +/// is carrying as many concurrent streams as it may - the lower of +/// [maxStreamsPerConnection] and the server's own advertised limit - a new +/// connection is dialed rather than queuing behind the existing one. +/// +/// Being multi-host makes this safe to use as a general-purpose transport - +/// for example as the `baseClient` passed to `googleapis_auth`'s client +/// helpers, whose credential negotiation (OAuth2 token endpoint, WIF/OIDC +/// token exchange) targets different hosts than the API calls that follow. +/// A connection dialed for one host is never reused for another. +/// +/// `onBadCertificate` is forwarded as-is to `SecureSocket.connect`: returning +/// `true` accepts a certificate that failed normal verification (expired, +/// self-signed, wrong host, ...). It exists for tests and trusted private +/// networks - do not use it to accept arbitrary certificates in production. +class Http2Client extends BaseClient { + Http2Client({ + this.maxStreamsPerConnection = 100, + this.maxIdleConnections = 1, + this.settingsTimeout = const Duration(seconds: 10), + int maxConcurrentHandshakes = 50, + SecurityContext? context, + bool Function(X509Certificate certificate)? onBadCertificate, + }) : _context = context, + _onBadCertificate = onBadCertificate, + _handshakeGate = Pool(maxConcurrentHandshakes); + + /// The maximum number of concurrent HTTP/2 streams (i.e. requests) to + /// multiplex onto a single connection before dialing another. + /// + /// An upper bound only: if a server says it accepts fewer concurrent + /// streams than this, that smaller number is used for its connections. + final int maxStreamsPerConnection; + + /// The maximum number of idle connections to keep per host, per + /// [ClientPool.maxIdleResources]. + final int maxIdleConnections; + + /// How long to wait for a freshly dialed connection's peer to send its + /// mandatory initial SETTINGS frame (RFC 7540 3.5). + /// + /// Until it arrives the peer's stream limit is unknown, so the connection + /// can't be safely multiplexed onto; a peer that never sends one is + /// abandoned rather than waited on forever. + final Duration settingsTimeout; + + final SecurityContext? _context; + + // Forwarded to `SecureSocket.connect` as-is: returning `true` accepts a + // certificate that failed normal verification (expired, self-signed, + // wrong host, ...). Intended for tests and trusted private networks - + // never use this to accept arbitrary certificates in production. + final bool Function(X509Certificate certificate)? _onBadCertificate; + + // Caps concurrent in-flight TCP+TLS handshakes across every host, + // independent of how many connections any one host's pool ends up + // needing. + final Pool _handshakeGate; + + final _pools = >{}; + var _closed = false; + + // Synchronous (no `await`), so concurrent requests to a new host:port + // can't race each other into creating two pools for the same key. + ClientPool _poolFor(Uri url) { + final key = '${url.host}:${url.port}'; + return _pools.putIfAbsent( + key, + () => ClientPool( + () => _handshakeGate.withResource(() => _dial(url.host, url.port)), + maxConcurrentOperations: maxStreamsPerConnection, + maxIdleResources: maxIdleConnections, + destroy: (transport) => transport.finish(), + concurrencyLimitOf: (transport) => transport.peerMaxConcurrentStreams, + ), + ); + } + + Future _dial(String host, int port) async { + final socket = await SecureSocket.connect( + host, + port, + context: _context, + onBadCertificate: _onBadCertificate, + supportedProtocols: ['h2'], + ); + if (socket.selectedProtocol != 'h2') { + socket.destroy(); + throw StateError( + 'Server did not negotiate HTTP/2 (got ${socket.selectedProtocol})', + ); + } + + final died = Completer(); + final incoming = socket.transform( + StreamTransformer>.fromHandlers( + handleDone: (sink) { + if (!died.isCompleted) died.complete(); + sink.close(); + }, + handleError: (error, stackTrace, sink) { + if (!died.isCompleted) died.completeError(error, stackTrace); + sink.addError(error, stackTrace); + }, + ), + ); + unawaited(died.future.catchError((Object _) {})); + + final transport = ClientConnection( + incoming, + socket, + const ClientSettings(), + ); + try { + await Future.any([ + transport.onInitialPeerSettingsReceived, + died.future.then( + (_) => + throw ClientException( + 'The connection closed before the peer sent its initial ' + 'SETTINGS frame.', + ), + ), + ]).timeout(settingsTimeout); + } catch (_) { + await transport.terminate(); + rethrow; + } + return transport; + } + + /// Sends [request] as a single HTTP/2 stream on [lease]'s connection, + /// translating between `BaseRequest`/`StreamedResponse` and http2's frames. + /// + /// Returns once the response headers arrive, but takes ownership of [lease]: + /// the stream stays open while the body is delivered, so the slot is only + /// released once that stream reaches a terminal state. + Future _sendOverHttp2( + PoolLease lease, + BaseRequest request, + List bodyBytes, + ) async { + final transport = lease.value; + if (transport.isClosing) throw const _ConnectionClosedByPeer(); + if (!transport.canOpenStream) throw const _ConnectionAtStreamLimit(); + + final rawPath = request.url.path.isEmpty ? '/' : request.url.path; + final path = + request.url.hasQuery ? '$rawPath?${request.url.query}' : rawPath; + + final stream = transport.makeRequest([ + Header.ascii(':method', request.method), + Header.ascii(':scheme', 'https'), + Header.ascii(':authority', _authorityOf(request.url)), + Header.ascii(':path', path), + for (final entry in request.headers.entries) + if (!_connectionSpecificHeaders.contains(entry.key.toLowerCase())) + Header.ascii(entry.key.toLowerCase(), entry.value), + ], endStream: bodyBytes.isEmpty); + + if (bodyBytes.isNotEmpty) stream.sendData(bodyBytes, endStream: true); + + final statusCompleter = Completer(); + late final StreamSubscription subscription; + final bodyController = StreamController>( + onCancel: () { + lease.release(); + stream.terminate(); + return subscription.cancel(); + }, + ); + final responseHeaders = {}; + + subscription = stream.incomingMessages.listen( + (message) { + if (message is HeadersStreamMessage) { + for (final header in message.headers) { + final name = latin1.decode(header.name); + final value = latin1.decode(header.value); + if (name == ':status') { + final status = int.tryParse(value); + if (status == null) { + if (!statusCompleter.isCompleted) { + statusCompleter.completeError( + ClientException( + 'Invalid HTTP/2 ":status" value "$value"', + request.url, + ), + ); + } + } else if (status >= 200 && !statusCompleter.isCompleted) { + statusCompleter.complete(status); + } + } else { + responseHeaders.update( + name, + (existing) => '$existing, $value', + ifAbsent: () => value, + ); + } + } + } else if (message is DataStreamMessage) { + bodyController.add(message.bytes); + } + }, + onDone: () { + if (!statusCompleter.isCompleted) { + statusCompleter.completeError( + ClientException( + 'Stream closed before a response status was received', + request.url, + ), + ); + } + if (!bodyController.isClosed) { + bodyController.close(); + } + lease.release(); + }, + onError: (Object error, StackTrace stackTrace) { + final failure = + error is ClientException + ? error + : ClientException('$error', request.url); + if (!statusCompleter.isCompleted) { + statusCompleter.completeError(failure, stackTrace); + } + bodyController.addError(failure, stackTrace); + if (!bodyController.isClosed) { + bodyController.close(); + } + lease.release(); + }, + cancelOnError: true, + ); + + final statusCode = await statusCompleter.future; + return StreamedResponse( + bodyController.stream, + statusCode, + contentLength: int.tryParse(responseHeaders['content-length'] ?? ''), + headers: Map.unmodifiable(responseHeaders), + reasonPhrase: _reasonPhrases[statusCode], + request: request, + ); + } + + /// The number of connections currently pooled, across every host. + int get connectionCount => _pools.values.map((pool) => pool.size).sum; + + @override + Future send(BaseRequest request) async { + if (_closed) { + throw ClientException( + 'HTTP request failed. Client is already closed.', + request.url, + ); + } + if (request.url.scheme != 'https') { + throw ClientException( + 'Http2Client only supports https (got "${request.url.scheme}").', + request.url, + ); + } + + List? bodyBytes; + + Future attempt() async { + final lease = await _poolFor(request.url).acquire(); + try { + bodyBytes ??= await request.finalize().toBytes(); + } catch (_) { + lease.release(); + rethrow; + } + try { + return await _sendOverHttp2(lease, request, bodyBytes!); + } on _ConnectionAtStreamLimit { + // Busy, not broken: give the slot back but leave the connection in + // the pool for whoever needs it next. + lease.release(); + rethrow; + } catch (_) { + lease.markFailed(); + lease.release(); + rethrow; + } + } + + return attempt() + .catchError( + (Object _) => attempt(), + test: + (error) => + error is _ConnectionClosedByPeer || + error is _ConnectionAtStreamLimit, + ) + .catchError( + (Object error, StackTrace stackTrace) => Error.throwWithStackTrace( + ClientException('$error', request.url), + stackTrace, + ), + test: (error) => error is! ClientException, + ); + } + + /// Waits for in-flight requests to finish, then closes every connection. + /// + /// Unlike [close] (constrained by `http.Client`'s synchronous signature), + /// this can be awaited by callers who hold a concrete [Http2Client]. + /// + /// A request counts as in-flight until its response body ends or is + /// cancelled, so a caller holding a response it never reads will hold this + /// up. [close] does not await this, so it can never block on that. + Future terminate() async { + _closed = true; + final pools = _pools.values.toList(); + _pools.clear(); + await Future.wait(pools.map((pool) => pool.terminate())); + } + + @override + void close() => unawaited(terminate()); +} + +/// Thrown by [Http2Client._sendOverHttp2] when a pooled connection turns +/// out to have already been closed by the peer (e.g. a graceful `GOAWAY`) +/// before any bytes were written for this request. [Http2Client.send] +/// catches this and retries once on whatever the pool dials next - +/// [ClientPool.run] has already marked the dead connection failed and +/// evicted it by the time the retry runs. +class _ConnectionClosedByPeer implements Exception { + const _ConnectionClosedByPeer(); + + @override + String toString() => + 'The pooled HTTP/2 connection was closed by the peer before this ' + 'request could be sent.'; +} + +/// Thrown by [Http2Client._sendOverHttp2] when a pooled connection is alive +/// but already carrying as many concurrent streams as its peer allows. +/// +/// Unlike [_ConnectionClosedByPeer] this says nothing bad about the +/// connection, so [Http2Client.send] retries without condemning it. +class _ConnectionAtStreamLimit implements Exception { + const _ConnectionAtStreamLimit(); + + @override + String toString() => + 'The pooled HTTP/2 connection is already carrying as many concurrent ' + 'streams as the server allows.'; +} + +/// HTTP/2 carries no reason phrase (RFC 9113 8.3.2 dropped it as redundant +/// with the status code), so one is derived from the status instead - the same +/// approach `package:cupertino_http` takes for NSURLSession. +const _reasonPhrases = { + 100: 'Continue', + 101: 'Switching Protocols', + 200: 'OK', + 201: 'Created', + 202: 'Accepted', + 203: 'Non-Authoritative Information', + 204: 'No Content', + 205: 'Reset Content', + 206: 'Partial Content', + 300: 'Multiple Choices', + 301: 'Moved Permanently', + 302: 'Found', + 303: 'See Other', + 304: 'Not Modified', + 305: 'Use Proxy', + 307: 'Temporary Redirect', + 308: 'Permanent Redirect', + 400: 'Bad Request', + 401: 'Unauthorized', + 402: 'Payment Required', + 403: 'Forbidden', + 404: 'Not Found', + 405: 'Method Not Allowed', + 406: 'Not Acceptable', + 407: 'Proxy Authentication Required', + 408: 'Request Time-out', + 409: 'Conflict', + 410: 'Gone', + 411: 'Length Required', + 412: 'Precondition Failed', + 413: 'Request Entity Too Large', + 414: 'Request-URI Too Long', + 415: 'Unsupported Media Type', + 416: 'Requested range not satisfiable', + 417: 'Expectation Failed', + 421: 'Misdirected Request', + 422: 'Unprocessable Entity', + 426: 'Upgrade Required', + 428: 'Precondition Required', + 429: 'Too Many Requests', + 431: 'Request Header Fields Too Large', + 500: 'Internal Server Error', + 501: 'Not Implemented', + 502: 'Bad Gateway', + 503: 'Service Unavailable', + 504: 'Gateway Time-out', + 505: 'Http Version not supported', + 511: 'Network Authentication Required', +}; + +/// RFC 9113 8.3.1: ":authority" carries the port unless it is the default for +/// the scheme - which is always https here, so 443. +String _authorityOf(Uri url) => + url.port == 443 ? url.host : '${url.host}:${url.port}'; + +/// Header fields forbidden on an HTTP/2 stream (RFC 7540 8.1.2.2), plus +/// `host` since `:authority` already carries what it would. +const _connectionSpecificHeaders = { + 'connection', + 'keep-alive', + 'proxy-connection', + 'transfer-encoding', + 'upgrade', + 'host', +}; diff --git a/pkgs/http2/pubspec.yaml b/pkgs/http2/pubspec.yaml index e65f0a31a4..a14707ec97 100644 --- a/pkgs/http2/pubspec.yaml +++ b/pkgs/http2/pubspec.yaml @@ -11,8 +11,20 @@ topics: environment: sdk: ^3.7.0 +dependencies: + collection: ^1.19.0 + http: ^1.5.0 + meta: ^1.15.0 + pool: ^1.5.0 + dev_dependencies: build_runner: ^2.4.15 dart_flutter_team_lints: ^3.5.1 + http_client_conformance_tests: + path: ../http_client_conformance_tests/ mockito: ^5.4.5 test: ^1.25.15 + +dependency_overrides: + http: + path: ../http/ diff --git a/pkgs/http2/test/client_conformance_test.dart b/pkgs/http2/test/client_conformance_test.dart new file mode 100644 index 0000000000..0294ffe77c --- /dev/null +++ b/pkgs/http2/test/client_conformance_test.dart @@ -0,0 +1,283 @@ +// Copyright (c) 2026, the Dart project authors. Please see the AUTHORS file +// for details. All rights reserved. Use of this source code is governed by a +// BSD-style license that can be found in the LICENSE file. + +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; + +import 'package:http/http.dart' as http; +import 'package:http/io_client.dart'; +import 'package:http2/src/http2_client.dart'; +import 'package:http2/transport.dart'; +import 'package:http_client_conformance_tests/http_client_conformance_tests.dart'; +import 'package:test/test.dart'; + +class Http2ProxyServer { + final SecureServerSocket _socket; + final List _connections = []; + final IOClient _httpClient = IOClient(); + + Http2ProxyServer._(this._socket) { + _socket.listen((socket) { + final connection = ServerTransportConnection.viaSocket(socket); + _connections.add(connection); + connection.incomingStreams.listen(_handleStream); + }); + } + + static Future start() async { + final context = + SecurityContext() + ..useCertificateChain('test/certificates/server_chain.pem') + ..usePrivateKey( + 'test/certificates/server_key.pem', + password: 'dartdart', + ) + ..setAlpnProtocols(['h2'], true); + final socket = await SecureServerSocket.bind('localhost', 0, context); + return Http2ProxyServer._(socket); + } + + int get port => _socket.port; + + Future _handleStream(ServerTransportStream stream) async { + try { + final messages = StreamIterator(stream.incomingMessages); + if (!await messages.moveNext()) return; + + final headersMsg = messages.current as HeadersStreamMessage; + String? method; + String? path; + int? targetPort; + final headers = {}; + + for (final header in headersMsg.headers) { + final name = ascii.decode(header.name); + final value = ascii.decode(header.value); + if (name == ':method') { + method = value; + } else if (name == ':path') { + path = value; + } else if (name == 'x-target-port') { + targetPort = int.parse(value); + } else if (!name.startsWith(':')) { + headers[name] = value; + } + } + + if (method == null || path == null || targetPort == null) { + stream.outgoingMessages.add( + HeadersStreamMessage([Header.ascii(':status', '400')]), + ); + await stream.outgoingMessages.close(); + return; + } + + // Collect body + final bodyBytes = []; + while (await messages.moveNext()) { + final msg = messages.current; + if (msg is DataStreamMessage) { + bodyBytes.addAll(msg.bytes); + } + } + + // Forward to HTTP/1.1 server + final targetUri = Uri.parse('http://localhost:$targetPort$path'); + final httpRequest = http.Request(method, targetUri); + headers.forEach((k, v) { + httpRequest.headers[k] = v; + }); + httpRequest.bodyBytes = bodyBytes; + + final httpResponse = await _httpClient.send(httpRequest); + + // Send response headers + final responseHeaders =
[ + Header.ascii(':status', httpResponse.statusCode.toString()), + ]; + httpResponse.headers.forEach((k, v) { + responseHeaders.add(Header.ascii(k.toLowerCase(), v)); + }); + stream.outgoingMessages.add(HeadersStreamMessage(responseHeaders)); + + // Send response body + await for (final chunk in httpResponse.stream) { + stream.outgoingMessages.add(DataStreamMessage(chunk)); + } + await stream.outgoingMessages.close(); + } catch (e) { + print('Proxy error: $e'); + stream.terminate(); + } + } + + Future close() async { + await _socket.close(); + for (final conn in _connections) { + await conn.terminate(); + } + _httpClient.close(); + } +} + +class ProxyRequest extends http.BaseRequest implements http.Abortable { + // `url` is not redeclared here - BaseRequest's own constructor stores it. + ProxyRequest(this._original, Uri url) : super(_original.method, url); + + final http.BaseRequest _original; + + @override + Future? get abortTrigger => + _original is http.Abortable ? _original.abortTrigger : null; + + @override + Map get headers => _original.headers; + + @override + int? get contentLength => _original.contentLength; + @override + set contentLength(int? value) => _original.contentLength = value; + + @override + bool get followRedirects => _original.followRedirects; + @override + set followRedirects(bool value) => _original.followRedirects = value; + + @override + int get maxRedirects => _original.maxRedirects; + @override + set maxRedirects(int value) => _original.maxRedirects = value; + + @override + bool get persistentConnection => _original.persistentConnection; + @override + set persistentConnection(bool value) => + _original.persistentConnection = value; + + @override + http.ByteStream finalize() { + super.finalize(); + return _original.finalize(); + } +} + +class ConformanceProxyClient extends http.BaseClient { + final Http2Client _inner; + final int _proxyPort; + + ConformanceProxyClient(this._inner, this._proxyPort); + + @override + Future send(http.BaseRequest request) async { + final targetPort = request.url.port; + final proxyUrl = request.url.replace( + scheme: 'https', + host: 'localhost', + port: _proxyPort, + ); + + request.headers['x-target-port'] = targetPort.toString(); + final proxyRequest = ProxyRequest(request, proxyUrl); + + try { + final response = await _inner.send(proxyRequest); + return http.StreamedResponse( + response.stream, + response.statusCode, + contentLength: response.contentLength, + headers: response.headers, + isRedirect: response.isRedirect, + persistentConnection: response.persistentConnection, + reasonPhrase: response.reasonPhrase, + request: request, + ); + } on http.ClientException catch (e) { + if (e is http.RequestAbortedException) { + throw http.RequestAbortedException(request.url); + } + throw http.ClientException(e.message, request.url); + } + } + + @override + void close() { + _inner.close(); + } +} + +/// [Http2Client] only supports HTTP/2 over TLS (HTTPS). However, the standard +/// servers started by http_client_conformance_tests only support unencrypted +/// HTTP/1.1. +/// To bridge this protocol gap, we run a local HTTP/2 proxy server +/// ([Http2ProxyServer]) in-process. [ConformanceProxyClient] wraps +/// [Http2Client] and rewrites the destination URI of all outgoing requests to +/// point to the local proxy server, attaching a custom `x-target-port` header +/// to specify the target HTTP/1.1 server port. The proxy server then forwards +/// the request over HTTP/1.1 and returns the response to [Http2Client] over +/// HTTP/2. +void main() { + late final Http2ProxyServer proxy; + + setUpAll(() async { + proxy = await Http2ProxyServer.start(); + }); + + tearDownAll(() async { + await proxy.close(); + }); + + ConformanceProxyClient clientFactory() => ConformanceProxyClient( + Http2Client(onBadCertificate: (_) => true), + proxy.port, + ); + + testRequestBody(clientFactory); + + // TODO: Implement request body streaming support in Http2Client. + // Currently Http2Client reads the entire request body into memory before + // sending. + testRequestBodyStreamed(clientFactory, canStreamRequestBody: false); + + testResponseBody(clientFactory); + // TODO: Re-enable once request abort support is implemented in Http2Client. + // testResponseBodyStreamed(clientFactory); + testRequestHeaders(clientFactory); + testRequestMethods(clientFactory, preservesMethodCase: false); + + testResponseHeaders( + clientFactory, + // HTTP/2 explicitly forbids folded headers (RFC 7540 Section 8.1.2.6). + supportsFoldedHeaders: false, + // HTTP/2 does not allow NUL characters inside header names or values. + correctlyHandlesNullHeaderValues: false, + ); + + testResponseStatusLine(clientFactory); + + // TODO: Implement redirect-following support in Http2Client. + // testRedirect(clientFactory); + + testServerErrors(clientFactory); + testCompressedResponseBody(clientFactory); + testMultipleClients(clientFactory); + testMultipartRequests(clientFactory, supportsMultipartRequest: true); + testClose(clientFactory); + + // TODO: Support running client conformance tests in isolates. + // Currently we set `canWorkInIsolates` to false because the proxy server uses + // `SecureServerSocket`, which cannot be sent across isolates. + testIsolate(clientFactory, canWorkInIsolates: false); + + testRequestCookies(clientFactory, canSendCookieHeaders: true); + testResponseCookies(clientFactory, canReceiveSetCookieHeaders: true); + + // TODO: Implement request abort support in Http2Client. + // testAbort( + // clientFactory, + // supportsAbort: true, + // canStreamRequestBody: false, + // canStreamResponseBody: true, + // ); +} diff --git a/pkgs/http2/test/client_pool_test.dart b/pkgs/http2/test/client_pool_test.dart new file mode 100644 index 0000000000..a970b7ac10 --- /dev/null +++ b/pkgs/http2/test/client_pool_test.dart @@ -0,0 +1,465 @@ +// Copyright (c) 2026, the Dart project authors. Please see the AUTHORS file +// for details. All rights reserved. Use of this source code is governed by a +// BSD-style license that can be found in the LICENSE file. + +import 'dart:async'; +import 'dart:math'; + +import 'package:http2/src/client_pool.dart'; +import 'package:test/test.dart'; + +ClientPool _pool({ + required int maxConcurrentOperations, + int maxIdleResources = 1, + int? Function(int resource)? concurrencyLimitOf, + List? destroyed, + Future? destroyGate, + Future? dialGate, +}) { + var nextId = 0; + return ClientPool( + () async { + final id = nextId++; + if (dialGate != null) await dialGate; + return id; + }, + maxConcurrentOperations: maxConcurrentOperations, + maxIdleResources: maxIdleResources, + destroy: (resource) async { + destroyed?.add(resource); + if (destroyGate != null) await destroyGate; + }, + concurrencyLimitOf: concurrencyLimitOf, + ); +} + +void main() { + group('client-pool-test', () { + test('creates-new-resources-as-needed', () { + final pool = _pool(maxConcurrentOperations: 2); + final completers = List.generate(3, (_) => Completer()); + + expect(pool.size, 0); + unawaited(pool.run((_) => completers[0].future)); + unawaited(pool.run((_) => completers[1].future)); + expect(pool.size, 1); + unawaited(pool.run((_) => completers[2].future)); + expect(pool.size, 2); + + for (final c in completers) { + c.complete(); + } + }); + + test('reuses-resource-with-remaining-capacity', () async { + final pool = _pool(maxConcurrentOperations: 2); + final completers = List.generate(3, (_) => Completer()); + + final first = pool.run((_) => completers[0].future); + unawaited(pool.run((_) => completers[1].future)); + expect(pool.size, 1); + + completers[0].complete(); + await first; + + unawaited(pool.run((_) => completers[2].future)); + expect(pool.size, 1); + + completers[1].complete(); + completers[2].complete(); + }); + + test('packs-load-onto-most-full-resource', () async { + final pool = _pool(maxConcurrentOperations: 2); + final completers = List.generate(4, (_) => Completer()); + final resourcesUsed = []; + + void run(int i) => unawaited( + pool.run((r) { + resourcesUsed.add(r); + return completers[i].future; + }), + ); + + run(0); + run(1); + run(2); // Resource 0 is full - this should open resource 1. + await Future.value(); + expect(resourcesUsed, [0, 0, 1]); + + completers[0].complete(); + await Future.value(); + + run(3); // Resource 0 has a free slot again and is the most-full option. + await Future.value(); + expect(resourcesUsed, [0, 0, 1, 0]); + + completers[1].complete(); + completers[2].complete(); + completers[3].complete(); + }); + + test('stops-reusing-resource-after-failure', () async { + final pool = _pool(maxConcurrentOperations: 10); + final resourcesUsed = []; + + await pool + .run((r) { + resourcesUsed.add(r); + return Future.error('boom'); + }) + .catchError((_) {}); + + await pool.run((r) { + resourcesUsed.add(r); + return Future.value(); + }); + + expect(resourcesUsed, [0, 1]); + }); + + test('garbage-collects-after-success', () async { + final pool = _pool(maxConcurrentOperations: 2, maxIdleResources: 0); + final completers = List.generate(4, (_) => Completer()); + + final ops = [ + pool.run((_) => completers[0].future), + pool.run((_) => completers[1].future), + pool.run((_) => completers[2].future), + pool.run((_) => completers[3].future), + ]; + expect(pool.size, 2); + + for (final c in completers) { + c.complete(); + } + await Future.wait(ops); + + expect(pool.size, 0); + }); + + test('garbage-collects-after-error', () async { + final pool = _pool(maxConcurrentOperations: 2, maxIdleResources: 0); + + final ops = List.generate( + 4, + (_) => pool.run((_) => Future.error('boom')).catchError((_) {}), + ); + await Future.wait(ops); + + expect(pool.size, 0); + }); + + test('keeps-idle-resources-up-to-max-idle-resources', () async { + final pool = _pool(maxConcurrentOperations: 1, maxIdleResources: 3); + final completers = List.generate(4, (_) => Completer()); + + final ops = [ + pool.run((_) => completers[0].future), + pool.run((_) => completers[1].future), + pool.run((_) => completers[2].future), + pool.run((_) => completers[3].future), + ]; + expect(pool.size, 4); + + for (final c in completers) { + c.complete(); + } + await Future.wait(ops); + + expect(pool.size, 3); + }); + + test('honours-a-resources-own-lower-concurrency-limit', () async { + final pool = _pool( + maxConcurrentOperations: 10, + concurrencyLimitOf: (_) => 2, + ); + final completers = List.generate(3, (_) => Completer()); + final resourcesUsed = []; + + void run(int i) => unawaited( + pool.run((r) { + resourcesUsed.add(r); + return completers[i].future; + }), + ); + + run(0); + await pool.run((_) async {}); + run(1); + run(2); // Resource 0 is at its own limit of 2 - this opens resource 1. + await Future.value(); + + expect(resourcesUsed, [0, 0, 1]); + expect(pool.size, 2); + + for (final c in completers) { + c.complete(); + } + }); + + test('ignores-a-resource-limit-above-max-concurrent-operations', () async { + final pool = _pool( + maxConcurrentOperations: 1, + concurrencyLimitOf: (_) => 1000, + ); + final completers = List.generate(2, (_) => Completer()); + + unawaited(pool.run((_) => completers[0].future)); + await Future.value(); + unawaited(pool.run((_) => completers[1].future)); + await Future.value(); + + expect(pool.size, 2); + + for (final c in completers) { + c.complete(); + } + }); + + test('re-reads-a-resource-limit-that-changes', () async { + var limit = 2; + final pool = _pool( + maxConcurrentOperations: 10, + concurrencyLimitOf: (_) => limit, + ); + final completers = List.generate(3, (_) => Completer()); + final resourcesUsed = []; + + void run(int i) => unawaited( + pool.run((r) { + resourcesUsed.add(r); + return completers[i].future; + }), + ); + + run(0); + await pool.run((_) async {}); // Let resource 0 resolve. + run(1); // Still within the limit of 2, so resource 0 is reused. + await Future.value(); + expect(resourcesUsed, [0, 0]); + + limit = 1; + run(2); // Resource 0 is now over its lowered limit, so a new one opens. + await Future.value(); + expect(resourcesUsed, [0, 0, 1]); + + for (final c in completers) { + c.complete(); + } + }); + + test('keeps-a-limited-resource-that-is-merely-full', () async { + final destroyed = []; + final pool = _pool( + maxConcurrentOperations: 10, + concurrencyLimitOf: (_) => 1, + destroyed: destroyed, + ); + final completers = List.generate(2, (_) => Completer()); + + final first = pool.run((_) => completers[0].future); + await Future.value(); + final second = pool.run((_) => completers[1].future); + await Future.value(); + expect(pool.size, 2); + + completers[0].complete(); + await first; + expect(destroyed, isEmpty); + expect(pool.size, 2); + + completers[1].complete(); + await second; + + expect(destroyed, hasLength(1)); + expect(pool.size, 1); + }); + + test('rejects-operations-after-terminate', () async { + final pool = _pool(maxConcurrentOperations: 1); + + await pool.terminate(); + + expect(() => pool.run((_) async {}), throwsA(isA())); + }); + + test( + 'does-not-over-commit-a-resource-whose-limit-is-not-known-yet', + () async { + final dialGate = Completer(); + final workGate = Completer(); + final pool = _pool( + maxConcurrentOperations: 100, + maxIdleResources: 10, + concurrencyLimitOf: (_) => 1, + dialGate: dialGate.future, + ); + + final inFlight = {}; + final peak = {}; + final ops = List.generate( + 8, + (_) => pool.run((r) async { + final now = (inFlight[r] ?? 0) + 1; + inFlight[r] = now; + peak[r] = max(peak[r] ?? 0, now); + await workGate.future; + inFlight[r] = inFlight[r]! - 1; + }), + ); + + dialGate.complete(); + await pumpEventQueue(); + + expect( + peak.values, + everyElement(1), + reason: 'no resource may run more than its own limit of 1', + ); + expect(peak, hasLength(8)); + + workGate.complete(); + await Future.wait(ops); + }, + ); + + test('tolerates-a-resource-that-reports-a-zero-limit', () async { + final pool = _pool( + maxConcurrentOperations: 10, + concurrencyLimitOf: (_) => 0, + ); + + await expectLater(pool.run((r) async => r), completion(0)); + }); + + test('a-held-lease-keeps-its-slot', () async { + final pool = _pool(maxConcurrentOperations: 2); + + final lease = await pool.acquire(); + expect(pool.opCount, 1); + + lease.release(); + expect(pool.opCount, 0); + }); + + test('releasing-a-lease-twice-is-a-no-op', () async { + final pool = _pool(maxConcurrentOperations: 2); + + final lease = await pool.acquire(); + lease.release(); + lease.release(); + + expect(pool.opCount, 0); + }); + + test('a-failed-lease-stops-the-resource-being-reused', () async { + final pool = _pool(maxConcurrentOperations: 10); + final resourcesUsed = []; + + final lease = await pool.acquire(); + resourcesUsed.add(lease.value); + lease.markFailed(); + lease.release(); + + await pool.run((r) async => resourcesUsed.add(r)); + + expect(resourcesUsed, [0, 1]); + }); + + test('does-not-couple-an-operation-to-a-slow-destroy', () async { + final pool = _pool( + maxConcurrentOperations: 1, + maxIdleResources: 0, + destroyGate: Completer().future, + ); + + await expectLater(pool.run((_) async => 'done'), completion('done')); + }); + + test('terminate-waits-for-a-destroy-started-by-collection', () async { + final gate = Completer(); + final destroyed = []; + final pool = _pool( + maxConcurrentOperations: 1, + maxIdleResources: 0, + destroyed: destroyed, + destroyGate: gate.future, + ); + + await pool.run((_) async {}); + expect(destroyed, [0]); + + var terminated = false; + final termination = pool.terminate().then((_) => terminated = true); + await Future.value(); + expect(terminated, isFalse, reason: 'destroy has not finished yet'); + + gate.complete(); + await termination; + expect(terminated, isTrue); + }); + + test('terminate-is-idempotent', () async { + final destroyed = []; + final pool = _pool(maxConcurrentOperations: 2, destroyed: destroyed); + final completer = Completer(); + + unawaited(pool.run((_) => completer.future)); + + final first = pool.terminate(); + final second = pool.terminate(); + completer.complete(); + await Future.wait([first, second]); + + expect(destroyed, [0]); // Destroyed once, not once per terminate() call. + expect(pool.size, 0); + }); + + test('terminate-twice-does-not-throw-concurrent-modification', () async { + final pool = _pool(maxConcurrentOperations: 1); + final completers = List.generate(2, (_) => Completer()); + + final ops = [ + pool.run((_) => completers[0].future), + pool.run((_) => completers[1].future), + ]; + await Future.value(); + expect(pool.size, 2); + + final terminations = [pool.terminate(), pool.terminate()]; + for (final c in completers) { + c.complete(); + } + await Future.wait(ops); + + await expectLater(Future.wait(terminations), completes); + }); + + test('terminate-after-terminate-returns-immediately', () async { + final destroyed = []; + final pool = _pool(maxConcurrentOperations: 1, destroyed: destroyed); + + await pool.run((_) async {}); + await pool.terminate(); + await pool.terminate(); + + expect(destroyed, hasLength(lessThanOrEqualTo(1))); + }); + + test('waits-for-in-flight-operations-before-terminating', () async { + final pool = _pool(maxConcurrentOperations: 1); + final completer = Completer(); + var terminated = false; + + unawaited(pool.run((_) => completer.future)); + final terminateOp = pool.terminate().then((_) => terminated = true); + + expect(terminated, isFalse); + completer.complete(); + await terminateOp; + expect(terminated, isTrue); + }); + }); +} diff --git a/pkgs/http2/test/http2_client_test.dart b/pkgs/http2/test/http2_client_test.dart new file mode 100644 index 0000000000..5c68517dad --- /dev/null +++ b/pkgs/http2/test/http2_client_test.dart @@ -0,0 +1,518 @@ +// Copyright (c) 2026, the Dart project authors. Please see the AUTHORS file +// for details. All rights reserved. Use of this source code is governed by a +// BSD-style license that can be found in the LICENSE file. + +import 'dart:async'; +import 'dart:convert' show ascii; +import 'dart:io'; +import 'dart:math'; + +import 'package:http/http.dart' show ClientException, Request; +import 'package:http2/multiprotocol_server.dart'; +import 'package:http2/src/connection.dart'; +import 'package:http2/src/http2_client.dart'; +import 'package:http2/transport.dart'; +import 'package:test/test.dart'; + +SecurityContext _serverContext() => + SecurityContext() + ..useCertificateChain('test/certificates/server_chain.pem') + ..usePrivateKey('test/certificates/server_key.pem', password: 'dartdart'); + +Future _bind() => + MultiProtocolHttpServer.bind('localhost', 0, _serverContext()); + +Http2Client _testClient({ + int maxStreamsPerConnection = 100, + int maxIdleConnections = 1, +}) => Http2Client( + maxStreamsPerConnection: maxStreamsPerConnection, + maxIdleConnections: maxIdleConnections, + onBadCertificate: (_) => true, +); + +/// A minimal HTTP/2-only server that (unlike [MultiProtocolHttpServer]) +/// exposes each accepted [ServerTransportConnection], so a test can finish +/// one connection gracefully while the server keeps listening for new ones. +class _RawHttp2Server { + _RawHttp2Server._( + this._socket, + this._settings, + this._responseDelay, + this._bodyGate, + ) { + _socket.listen((socket) { + final connection = ServerTransportConnection.viaSocket( + socket, + settings: _settings, + ); + connections.add(connection); + connection.incomingStreams.listen( + _respondWith('ok', delay: _responseDelay, bodyGate: _bodyGate), + ); + }); + } + + /// [settings] defaults to the same value `ServerTransportConnection` would + /// have applied on its own, so callers that don't care are unaffected. + static Future<_RawHttp2Server> bind({ + ServerSettings settings = const ServerSettings(concurrentStreamLimit: 1000), + Future? responseDelay, + Future? bodyGate, + }) async { + final context = _serverContext()..setAlpnProtocols(['h2'], true); + final socket = await SecureServerSocket.bind('localhost', 0, context); + return _RawHttp2Server._(socket, settings, responseDelay, bodyGate); + } + + final SecureServerSocket _socket; + final ServerSettings _settings; + final Future? _responseDelay; + final Future? _bodyGate; + final connections = []; + + int get port => _socket.port; + + Future close() async { + await _socket.close(); + for (final connection in connections) { + await connection.terminate(); + } + } +} + +/// Replies with [body] after waiting on [delay], if given. +/// +/// [bodyGate] holds the response open *after* its headers have been sent, so a +/// test can observe a request whose headers have arrived but whose stream is +/// still open. +void Function(ServerTransportStream) _respondWith( + String body, { + Future? delay, + Future? bodyGate, +}) { + return (stream) async { + final subscription = StreamIterator(stream.incomingMessages); + await subscription.moveNext(); // Consume the request headers. + while (await subscription.moveNext()) {} // Drain any request body. + + if (delay != null) await delay; + + stream.outgoingMessages.add( + HeadersStreamMessage([Header.ascii(':status', '200')]), + ); + if (bodyGate != null) await bodyGate; + try { + stream.outgoingMessages.add(DataStreamMessage(ascii.encode(body))); + await stream.outgoingMessages.close(); + } catch (_) {} + }; +} + +void main() { + group('http2-client-test', () { + test('sends-request-and-receives-response', () async { + final server = await _bind(); + server.startServing( + (request) {}, + expectAsync1(_respondWith('hello'), count: 1), + ); + + final client = _testClient(); + final response = await client.get( + Uri.parse('https://localhost:${server.port}/'), + ); + + expect(response.statusCode, 200); + expect(response.body, 'hello'); + + await client.terminate(); + await server.close(); + }); + + test('pools-connections-per-host-and-port', () async { + final serverA = await _bind(); + final serverB = await _bind(); + serverA.startServing( + (request) {}, + expectAsync1(_respondWith('a'), count: 1), + ); + serverB.startServing( + (request) {}, + expectAsync1(_respondWith('b'), count: 1), + ); + + final client = _testClient(); + await Future.wait([ + client.get(Uri.parse('https://localhost:${serverA.port}/')), + client.get(Uri.parse('https://localhost:${serverB.port}/')), + ]); + + expect(client.connectionCount, 2); + + await client.terminate(); + await Future.wait([serverA.close(), serverB.close()]); + }); + + test('exceeding-max-streams-per-connection-opens-new-connection', () async { + final server = await _bind(); + final releaseA = Completer(); + final releaseB = Completer(); + var requestNr = 0; + server.startServing( + (request) {}, + expectAsync1((stream) { + final release = requestNr++ == 0 ? releaseA : releaseB; + return _respondWith('r', delay: release.future)(stream); + }, count: 2), + ); + + final client = _testClient(maxStreamsPerConnection: 1); + final requestA = client.get( + Uri.parse('https://localhost:${server.port}/a'), + ); + await Future.delayed(const Duration(milliseconds: 50)); + final requestB = client.get( + Uri.parse('https://localhost:${server.port}/b'), + ); + + await Future.delayed(const Duration(milliseconds: 50)); + expect(client.connectionCount, 2); + + releaseA.complete(); + releaseB.complete(); + await Future.wait([requestA, requestB]); + + await client.terminate(); + await server.close(); + }); + + test('retries-once-when-pooled-connection-was-closed-by-peer', () async { + final server = await _RawHttp2Server.bind(); + final client = _testClient(); + + final r1 = await client.get( + Uri.parse('https://localhost:${server.port}/'), + ); + expect(r1.statusCode, 200); + expect(client.connectionCount, 1); + + await server.connections.single.finish(); + await Future.delayed(const Duration(milliseconds: 200)); + + final r2 = await client.get( + Uri.parse('https://localhost:${server.port}/'), + ); + expect(r2.statusCode, 200); + expect(r2.body, 'ok'); + + await client.terminate(); + await server.close(); + }); + + test('respects-server-advertised-max-concurrent-streams', () async { + final release = Completer(); + final server = await _RawHttp2Server.bind( + settings: const ServerSettings(concurrentStreamLimit: 1), + responseDelay: release.future, + ); + final client = _testClient( + maxStreamsPerConnection: 100, + maxIdleConnections: 5, + ); + + final requestA = client.get( + Uri.parse('https://localhost:${server.port}/a'), + ); + await Future.delayed(const Duration(milliseconds: 100)); + final requestB = client.get( + Uri.parse('https://localhost:${server.port}/b'), + ); + await Future.delayed(const Duration(milliseconds: 100)); + + release.complete(); + final responses = await Future.wait([requestA, requestB]); + expect(responses.map((r) => r.statusCode), everyElement(200)); + expect(server.connections, hasLength(2)); + + expect(client.connectionCount, 2); + + await client.terminate(); + await server.close(); + }); + + test('holds-a-pool-slot-until-the-response-body-completes', () async { + final gate = Completer(); + final server = await _bind(); + server.startServing( + (request) {}, + expectAsync1(_respondWith('ok', bodyGate: gate.future), count: 2), + ); + + final client = _testClient(maxStreamsPerConnection: 1); + final url = Uri.parse('https://localhost:${server.port}/'); + final first = await client.send(Request('GET', url)); + final second = await client.send(Request('GET', url)); + + expect(client.connectionCount, 2); + + gate.complete(); + expect(await first.stream.bytesToString(), 'ok'); + expect(await second.stream.bytesToString(), 'ok'); + + await client.terminate(); + await server.close(); + }); + + test('releases-the-slot-when-the-response-body-is-cancelled', () async { + final gate = Completer(); + final server = await _bind(); + var streamNr = 0; + server.startServing( + (request) {}, + expectAsync1((stream) { + final held = streamNr++ == 0 ? gate.future : null; + return _respondWith('ok', bodyGate: held)(stream); + }, count: 2), + ); + + final client = _testClient(maxStreamsPerConnection: 1); + final url = Uri.parse('https://localhost:${server.port}/'); + + final first = await client.send(Request('GET', url)); + await first.stream.listen((_) {}).cancel(); + + final second = await client.get(url); + expect(second.statusCode, 200); + expect(client.connectionCount, 1); + + gate.complete(); + await client.terminate(); + await server.close(); + }); + + test('releases-the-slot-when-the-response-body-errors', () async { + final gate = Completer(); + final server = await _RawHttp2Server.bind(bodyGate: gate.future); + final client = _testClient(maxStreamsPerConnection: 1); + final url = Uri.parse('https://localhost:${server.port}/'); + + final first = await client.send(Request('GET', url)); + await server.connections.single.terminate(); + await expectLater( + first.stream.drain(), + throwsA(isA()), + ); + + gate.complete(); + + final second = await client.get(url); + expect(second.statusCode, 200); + + await client.terminate(); + await server.close(); + }); + + test('does-not-exceed-the-server-stream-limit-on-a-cold-burst', () async { + const streamLimit = 2; + const requestCount = 12; + final release = Completer(); + final context = _serverContext()..setAlpnProtocols(['h2'], true); + final socket = await SecureServerSocket.bind('localhost', 0, context); + + final active = {}; + final peak = {}; + socket.listen((raw) { + final connection = ServerTransportConnection.viaSocket( + raw, + settings: const ServerSettings(concurrentStreamLimit: streamLimit), + ); + connection.incomingStreams.listen((stream) async { + final now = (active[connection] ?? 0) + 1; + active[connection] = now; + peak[connection] = max(peak[connection] ?? 0, now); + + final messages = StreamIterator(stream.incomingMessages); + await messages.moveNext(); + while (await messages.moveNext()) {} + await release.future; + stream.outgoingMessages.add( + HeadersStreamMessage([Header.ascii(':status', '200')]), + ); + stream.outgoingMessages.add(DataStreamMessage(ascii.encode('ok'))); + await stream.outgoingMessages.close(); + + active[connection] = active[connection]! - 1; + }); + }); + + final client = _testClient(maxStreamsPerConnection: 100); + final url = Uri.parse('https://localhost:${socket.port}/'); + final requests = List.generate(requestCount, (_) => client.get(url)); + + await pumpEventQueue(); + release.complete(); + final responses = await Future.wait(requests); + + expect(responses.map((r) => r.statusCode), everyElement(200)); + expect( + peak.values, + everyElement(lessThanOrEqualTo(streamLimit)), + reason: 'no connection may carry more streams than the server allows', + ); + + await client.terminate(); + await socket.close(); + }); + + test('fails-the-dial-when-the-peer-closes-before-settings', () async { + final context = _serverContext()..setAlpnProtocols(['h2'], true); + final socket = await SecureServerSocket.bind('localhost', 0, context); + socket.listen((connection) => connection.destroy()); + + final client = _testClient(); + await expectLater( + client.get(Uri.parse('https://localhost:${socket.port}/')), + throwsA(isA()), + ); + + await client.terminate(); + await socket.close(); + }); + + test('fails-the-dial-when-the-peer-never-sends-settings', () async { + final context = _serverContext()..setAlpnProtocols(['h2'], true); + final socket = await SecureServerSocket.bind('localhost', 0, context); + final held = []; + socket.listen(held.add); + + final client = Http2Client( + onBadCertificate: (_) => true, + settingsTimeout: const Duration(milliseconds: 200), + ); + await expectLater( + client.get(Uri.parse('https://localhost:${socket.port}/')), + throwsA(isA()), + ); + + await client.terminate(); + for (final connection in held) { + connection.destroy(); + } + await socket.close(); + }); + + test('separates-a-busy-connection-from-a-dead-one', () async { + // `isOpen` answers false for both, which is why Http2Client asks + // isClosing/canOpenStream instead. + final gate = Completer(); + final server = await _RawHttp2Server.bind( + settings: const ServerSettings(concurrentStreamLimit: 1), + bodyGate: gate.future, + ); + final socket = await SecureSocket.connect( + 'localhost', + server.port, + onBadCertificate: (_) => true, + supportedProtocols: ['h2'], + ); + final connection = ClientConnection( + socket, + socket, + const ClientSettings(), + ); + await connection.onInitialPeerSettingsReceived; + + expect(connection.isOpen, isTrue); + expect(connection.isClosing, isFalse); + expect(connection.canOpenStream, isTrue); + + connection + .makeRequest([ + Header.ascii(':method', 'GET'), + Header.ascii(':scheme', 'https'), + Header.ascii(':authority', 'localhost'), + Header.ascii(':path', '/'), + ], endStream: true) + .incomingMessages + // Terminating the connection below resets this stream. + .listen((_) {}, onError: (Object _) {}); + await pumpEventQueue(); + + // Occupying the server's only slot makes `isOpen` false even though the + // connection is perfectly healthy - the two new getters tell them apart. + expect(connection.isOpen, isFalse); + expect(connection.isClosing, isFalse, reason: 'busy, not dead'); + expect(connection.canOpenStream, isFalse); + + gate.complete(); + await connection.terminate(); + expect(connection.isClosing, isTrue, reason: 'now genuinely dead'); + + await server.close(); + }); + + test('keeps-a-connection-that-is-only-at-its-stream-limit', () async { + final gate = Completer(); + final server = await _RawHttp2Server.bind( + settings: const ServerSettings(concurrentStreamLimit: 1), + bodyGate: gate.future, + ); + // Idle connections are kept generously, so the final assertion reflects + // whether one was discarded as failed rather than collected as excess. + final client = _testClient( + maxStreamsPerConnection: 100, + maxIdleConnections: 5, + ); + final url = Uri.parse('https://localhost:${server.port}/'); + + final first = await client.send(Request('GET', url)); + expect(client.connectionCount, 1); + + final second = await client.send(Request('GET', url)); + expect(client.connectionCount, 2); + + gate.complete(); + expect(await first.stream.bytesToString(), 'ok'); + expect(await second.stream.bytesToString(), 'ok'); + + expect( + client.connectionCount, + 2, + reason: 'a connection that was merely busy must not be discarded', + ); + + await client.terminate(); + await server.close(); + }); + + test('terminate-waits-for-in-flight-request', () async { + final server = await _bind(); + final release = Completer(); + server.startServing( + (request) {}, + expectAsync1(_respondWith('done', delay: release.future), count: 1), + ); + + final client = _testClient(); + final request = client.get( + Uri.parse('https://localhost:${server.port}/'), + ); + + var terminated = false; + final terminateFuture = client.terminate().then((_) { + terminated = true; + }); + + await Future.delayed(const Duration(milliseconds: 50)); + expect(terminated, isFalse); + + release.complete(); + await request; + await terminateFuture; + expect(terminated, isTrue); + + await server.close(); + }); + }); +}