Skip to content
Open
Show file tree
Hide file tree
Changes from 7 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
1 change: 1 addition & 0 deletions packages/google_cloud_firestore/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
- Updated `Transaction.delete` and `Transaction.update` type constraints to accept `DocumentReference<Object?>`. (thanks to @Levin-Me)
- Made `Timestamp` encodable by adding `toJson` method. (thanks to @OutdatedGuy)
- Update dependency `googleapis_auth: ^2.3.3` to fix `auth/insufficient-permission` errors with Application Default Credentials that have no quota project set.
- Fixed slow and failing requests under high concurrency by pooling shared HTTP/2 connections by default instead of opening a new HTTP/1.1 connection per request.

## 0.5.2

Expand Down
270 changes: 264 additions & 6 deletions packages/google_cloud_firestore/lib/src/firestore_http_client.dart
Original file line number Diff line number Diff line change
Expand Up @@ -13,13 +13,17 @@
// limitations under the License.

import 'dart:async';
import 'dart:convert';
import 'dart:io';

import 'package:google_cloud/constants.dart' as google_cloud;
import 'package:google_cloud/google_cloud.dart' as google_cloud;
import 'package:google_cloud_firestore_v1/firestore.dart' as firestore_v1;
import 'package:googleapis_auth/auth_io.dart' as googleapis_auth;
import 'package:http/http.dart';
import 'package:http2/transport.dart' hide Settings;
import 'package:meta/meta.dart';
import 'package:pool/pool.dart';

import '../google_cloud_firestore.dart';
import 'environment.dart';
Expand Down Expand Up @@ -77,6 +81,248 @@ class EmulatorClient extends BaseClient implements googleapis_auth.AuthClient {
void close() => client.close();
}

/// Sends [request] as a single HTTP/2 stream on [transport], translating
/// between `BaseRequest`/`StreamedResponse` and http2's frames.
Future<StreamedResponse> _sendOverHttp2(
ClientTransportConnection transport,
BaseRequest request,
) async {
final bodyBytes = await request.finalize().toBytes();
final path = request.url.hasQuery
? '${request.url.path}?${request.url.query}'
: request.url.path;

final stream = transport.makeRequest([
Header.ascii(':method', request.method),
Header.ascii(':scheme', 'https'),
Header.ascii(':authority', request.url.host),
Header.ascii(':path', path),
for (final entry in request.headers.entries)
Header.ascii(entry.key.toLowerCase(), entry.value),
], endStream: bodyBytes.isEmpty);

if (bodyBytes.isNotEmpty) stream.sendData(bodyBytes, endStream: true);

final statusCompleter = Completer<int>();
final bodyController = StreamController<List<int>>();
final responseHeaders = <String, String>{};

stream.incomingMessages.listen(
Comment thread
demolaf marked this conversation as resolved.
Outdated
(message) {
if (message is HeadersStreamMessage) {
for (final header in message.headers) {
final name = ascii.decode(header.name);
final value = ascii.decode(header.value);
if (name == ':status') {
if (!statusCompleter.isCompleted) {
statusCompleter.complete(int.parse(value));
}
} else {
responseHeaders[name] = value;
}
}
} else if (message is DataStreamMessage) {
bodyController.add(message.bytes);
}
},
onDone: () {
if (!statusCompleter.isCompleted) {
statusCompleter.completeError(
StateError('Stream closed before a response status was received'),
);
}
if (!bodyController.isClosed) bodyController.close();
},
onError: (Object error, StackTrace stackTrace) {
if (!statusCompleter.isCompleted) {
statusCompleter.completeError(error, stackTrace);
}
bodyController.addError(error, stackTrace);
if (!bodyController.isClosed) bodyController.close();
},
cancelOnError: true,
);

final statusCode = await statusCompleter.future;
return StreamedResponse(
bodyController.stream,
statusCode,
headers: responseHeaders,
request: request,
);
}

class _PooledResource<T> {
_PooledResource(this.future);
final Future<T> future;
int inFlight = 0;
bool failed = false;
}

/// A pool of resources of type [T], mirroring nodejs-firestore's
/// `ClientPool` (dev/src/pool.ts): packs load onto the most-full resource
/// under [maxConcurrentOperations] instead of spreading evenly, opens a
/// new resource once existing ones are full, retires a resource once
/// [run] reports it failed, and garbage collects idle resources past
/// [maxIdleResources].
@internal
class ClientPool<T> {
ClientPool(
Future<T> Function() create, {
required this.maxConcurrentOperations,
required Future<void> Function(T resource) destroy,
this.maxIdleResources = 1,
}) : _create = create,
_destroy = destroy;

final Future<T> Function() _create;
final Future<void> Function(T resource) _destroy;
final int maxConcurrentOperations;
final int maxIdleResources;

final _resources = <_PooledResource<T>>[];
var _terminated = false;
Completer<void>? _drained;

@visibleForTesting
int get size => _resources.length;

@visibleForTesting
int get opCount =>
_resources.fold(0, (total, resource) => total + resource.inFlight);

/// Runs [operation] on an available (or newly created) resource.
Future<R> run<R>(Future<R> Function(T resource) operation) async {
if (_terminated) {
throw StateError('This pool has already been terminated.');
}

final pooled = _acquire();
pooled.inFlight++;
try {
return await operation(await pooled.future);
} catch (_) {
pooled.failed = true;
rethrow;
} finally {
pooled.inFlight--;
if (_terminated) {
_maybeCompleteDrain();
} else {
await _collectIfIdle(pooled);
}
Comment on lines +138 to +142

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

Awaiting _collectIfIdle inside the finally block blocks the completion of the run call (and thus the user's Firestore request) until the idle connection is fully destroyed. Since connection cleanup is a background task, we should run it unawaited to avoid adding unnecessary latency to the user's request.

Suggested change
if (_terminated) {
_maybeCompleteDrain();
} else {
await _collectIfIdle(pooled);
}
if (_terminated) {
_maybeCompleteDrain();
} else {
unawaited(_collectIfIdle(pooled));
}

}
}

// Synchronous (no `await`), so concurrent calls can't race each other
// into both creating a resource before either sees the other's.
_PooledResource<T> _acquire() {
_PooledResource<T>? selected;
for (final resource in _resources) {
if (resource.failed) continue;
if (resource.inFlight < maxConcurrentOperations &&
(selected == null || resource.inFlight > selected.inFlight)) {
selected = resource;
}
}
if (selected != null) return selected;

final resource = _PooledResource<T>(_create());
_resources.add(resource);
return resource;
}

Future<void> _collectIfIdle(_PooledResource<T> resource) async {
if (resource.inFlight > 0) return;
if (!resource.failed && !_hasExcessIdleCapacity) return;

_resources.remove(resource);
await _destroy(await resource.future);
}
Comment thread
demolaf marked this conversation as resolved.

bool get _hasExcessIdleCapacity {
final idleCapacity = _resources.fold(
0,
(total, resource) => total + (maxConcurrentOperations - resource.inFlight),
);
return idleCapacity > maxIdleResources * maxConcurrentOperations;
}

void _maybeCompleteDrain() {
final drained = _drained;
if (drained != null && !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.
Future<void> terminate() async {
_terminated = true;

if (opCount > 0) {
_drained = Completer<void>();
await _drained!.future;
}

for (final resource in _resources) {
try {
await _destroy(await resource.future);
} catch (_) {
// Best-effort: one resource failing to close shouldn't stop the
// rest from being destroyed.
}
}
_resources.clear();
}
}

/// Sends every request as a stream on a pool of shared HTTP/2 connections
/// to [host] - see [ClientPool]. Pure transport, no auth: wrapped with
/// googleapis_auth's client helpers in [FirestoreHttpClient._createClient]
/// for a refreshing [googleapis_auth.AuthClient].
class Http2ClientPool extends BaseClient {
Http2ClientPool(this.host, {this.maxStreamsPerConnection = 100})
: _pool = ClientPool<ClientTransportConnection>(
() => _handshakeGate.withResource(() => _dial(host)),
maxConcurrentOperations: maxStreamsPerConnection,
destroy: (transport) => transport.finish(),
);

final String host;
final int maxStreamsPerConnection;
final ClientPool<ClientTransportConnection> _pool;

// Caps concurrent in-flight TCP+TLS handshakes across all pools,
// independent of how many total connections any one pool needs.
static final _handshakeGate = Pool(50);

static Future<ClientTransportConnection> _dial(String host) async {
final socket = await SecureSocket.connect(
host,
443,
supportedProtocols: ['h2'],
);
if (socket.selectedProtocol != 'h2') {
socket.destroy();
throw StateError(
'Server did not negotiate HTTP/2 (got ${socket.selectedProtocol})',
);
}
return ClientTransportConnection.viaSocket(socket);
}

/// The number of connections currently in the pool.
int get connectionCount => _pool.size;

@override
Future<StreamedResponse> send(BaseRequest request) =>
_pool.run((transport) => _sendOverHttp2(transport, request));

@override
void close() => unawaited(_pool.terminate());
}
Comment thread
demolaf marked this conversation as resolved.

/// HTTP client wrapper for Firestore API operations.
///
/// Provides authenticated API access with automatic project ID discovery.
Expand Down Expand Up @@ -154,24 +400,35 @@ class FirestoreHttpClient {
/// Lazy-initialized HTTP client that's cached for reuse.
late final Future<googleapis_auth.AuthClient> _client = _createClient();

// googleapis_auth never closes a caller-supplied baseClient, so
// close() must reach this transport itself.
Http2ClientPool? _transport;

/// Creates the appropriate HTTP client based on emulator configuration.
Future<googleapis_auth.AuthClient> _createClient() async {
if (_isUsingEmulator) {
// Emulator: Create unauthenticated client.
// Emulator: plain HTTP/1.1, matching nodejs-firestore's REST-mode
// behavior (its gRPC mode uses HTTP/2, but this SDK is REST-only).
return EmulatorClient(Client());
}

// Production: Create authenticated client.
final transport = Http2ClientPool(
_settings.host ?? 'firestore.googleapis.com',
);
_transport = transport;

final serviceAccountCreds = credential.serviceAccountCredentials;
if (serviceAccountCreds != null) {
return googleapis_auth.clientViaServiceAccount(serviceAccountCreds, [
'https://www.googleapis.com/auth/cloud-platform',
]);
return googleapis_auth.clientViaServiceAccount(
serviceAccountCreds,
['https://www.googleapis.com/auth/cloud-platform'],
baseClient: transport,
);
}

// Fall back to Application Default Credentials
return googleapis_auth.clientViaApplicationDefaultCredentials(
scopes: ['https://www.googleapis.com/auth/cloud-platform'],
baseClient: transport,
);
}

Expand Down Expand Up @@ -204,5 +461,6 @@ class FirestoreHttpClient {
Future<void> close() async {
final client = await _client;
client.close();
_transport?.close();
}
Comment thread
demolaf marked this conversation as resolved.
}
2 changes: 2 additions & 0 deletions packages/google_cloud_firestore/pubspec.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,10 @@ dependencies:
google_cloud_type: ^0.5.2
googleapis_auth: ^2.3.3
http: ^1.6.0
http2: ^2.3.1
intl: ^0.20.0
meta: ^1.17.0
pool: ^1.5.1
retry: ^3.1.2

dev_dependencies:
Expand Down
Loading
Loading