Skip to content
Draft
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
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ export 'src/core/ref.dart'
QueryRef,
QueryResult,
DataSource;
export 'src/firebase_data_connect.dart';
export 'src/firebase_data_connect.dart' hide RoutingTransport;
export 'src/optional.dart'
show
Optional,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import 'package:firebase_core_platform_interface/firebase_core_platform_interfac
import 'package:firebase_data_connect/src/common/common_library.dart';
import 'package:firebase_data_connect/src/core/ref.dart';
import 'package:flutter/foundation.dart';
import 'package:meta/meta.dart';

import './network/rest_library.dart';
import './network/transport_library.dart';
Expand Down Expand Up @@ -116,7 +117,7 @@ class FirebaseDataConnect extends FirebasePlugin {
appCheck,
auth,
);
transport = _RoutingTransport(rest, ws);
transport = RoutingTransport(rest, ws);
}

@visibleForTesting
Expand Down Expand Up @@ -240,8 +241,12 @@ class FirebaseDataConnect extends FirebasePlugin {
}
}

class _RoutingTransport implements DataConnectTransport {
_RoutingTransport(this.rest, this.websocket);
/// @nodoc
@internal
// TODO: Move RoutingTransport to its own file under src/network/ to avoid exposing it publicly
// and hiding it in exports.
class RoutingTransport implements DataConnectTransport {
RoutingTransport(this.rest, this.websocket);
final RestTransport rest;
final WebSocketTransport websocket;

Expand Down Expand Up @@ -294,7 +299,7 @@ class _RoutingTransport implements DataConnectTransport {
Variables? vars,
String? token,
) {
if (websocket.isConnected) {
if (websocket.isConnected && websocket.hasActiveSubscriptions) {
return websocket.invokeMutation(
operationId, queryName, deserializer, serializer, vars, token);
}
Expand All @@ -311,7 +316,7 @@ class _RoutingTransport implements DataConnectTransport {
Variables? vars,
String? token,
) {
if (websocket.isConnected) {
if (websocket.isConnected && websocket.hasActiveSubscriptions) {
return websocket.invokeQuery(
operationId, queryName, deserializer, serialize, vars, token);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,8 @@ class WebSocketTransport implements DataConnectTransport {
this.sdkType,
this.appCheck, [
this.auth,
]) {
@visibleForTesting WebSocketChannel Function(Uri)? connectWebSocket,
]) : _connectWebSocket = connectWebSocket ?? WebSocketChannel.connect {
final protocol = (transportOptions.isSecure ?? true) ? 'wss' : 'ws';
final host = transportOptions.host;
final port = transportOptions.port ?? 443;
Expand Down Expand Up @@ -94,6 +95,7 @@ class WebSocketTransport implements DataConnectTransport {
});
}

final WebSocketChannel Function(Uri) _connectWebSocket;
FirebaseAuth? auth;
String? _currentUid;
// ignore: unused_field
Expand Down Expand Up @@ -184,16 +186,39 @@ class WebSocketTransport implements DataConnectTransport {
/// vetoes any reconnect for the operation that was still being set up.
int _pendingOperationSetups = 0;

static const Duration _idleDisconnectTimeout = Duration(seconds: 15);
Timer? _idleDisconnectTimer;

void _checkIdleAndDisconnect() {
if (_pendingOperationSetups > 0) return;
if (_pendingOperationSetups > 0) {
_cancelIdleDisconnectTimer();
return;
}
if (_streamListeners.isEmpty && _unaryListeners.isEmpty) {
_isExpectedDisconnect = true;
_disconnect();
_releaseWebSocketTransport();
_clearState();
_scheduleIdleDisconnect();
} else {
_cancelIdleDisconnectTimer();
}
}

void _scheduleIdleDisconnect() {
if (_idleDisconnectTimer != null) return; // Already scheduled
_idleDisconnectTimer = Timer(_idleDisconnectTimeout, () {
_idleDisconnectTimer = null;
if (_streamListeners.isEmpty && _unaryListeners.isEmpty && _pendingOperationSetups == 0) {
_isExpectedDisconnect = true;
_disconnect();
_releaseWebSocketTransport();
_clearState();
}
});
}

void _cancelIdleDisconnectTimer() {
_idleDisconnectTimer?.cancel();
_idleDisconnectTimer = null;
}

/// Re-establishes the connection on behalf of operations that are already
/// registered.
///
Expand Down Expand Up @@ -284,6 +309,9 @@ class WebSocketTransport implements DataConnectTransport {

bool get isConnected => _channel != null;

/// Returns true if there are active stream listeners.
bool get hasActiveSubscriptions => _streamListeners.isNotEmpty;

@visibleForTesting
Map<String, String> buildHeaders(String? authToken, String? appCheckToken) =>
_buildHeaders(authToken, appCheckToken);
Expand All @@ -307,6 +335,7 @@ class WebSocketTransport implements DataConnectTransport {
Future<void>? _connectionFuture;

Future<void> _ensureConnected(String? authToken) {
_cancelIdleDisconnectTimer();
if (_channel != null) return Future.value();
if (_connectionFuture != null) return _connectionFuture!;
_connectionFuture = _doConnect(authToken).whenComplete(() {
Expand All @@ -331,7 +360,7 @@ class WebSocketTransport implements DataConnectTransport {
// `done`/`error` from it cannot clobber the channel we are about to open.
unawaited(_channelSubscription?.cancel());

final channel = WebSocketChannel.connect(Uri.parse(_url));
final channel = _connectWebSocket(Uri.parse(_url));
_channel = channel;
_channelSubscription = channel.stream.listen(
_onMessage,
Expand Down Expand Up @@ -584,6 +613,7 @@ class WebSocketTransport implements DataConnectTransport {
// Ignore events from a socket that is no longer the active one.
if (!identical(_channel, channel)) return;
developer.log('WebSocket error: $error');
_cancelIdleDisconnectTimer();
_channel = null;
_isReconnecting = false;
_scheduleReconnect();
Expand All @@ -602,6 +632,7 @@ class WebSocketTransport implements DataConnectTransport {

void disconnect() {
_isExpectedDisconnect = true;
_cancelIdleDisconnectTimer();
_disconnect();
_releaseWebSocketTransport();
}
Expand All @@ -614,6 +645,7 @@ class WebSocketTransport implements DataConnectTransport {
// A `done` from a socket we already replaced must not null out the current
// channel: every later `_send` would be dropped on the floor silently.
if (!identical(_channel, channel)) return;
_cancelIdleDisconnectTimer();
_channel = null;
_isReconnecting = false;
if (!_isExpectedDisconnect) {
Expand Down Expand Up @@ -682,6 +714,7 @@ class WebSocketTransport implements DataConnectTransport {
authToken, requestKind, isMutation);
} finally {
_pendingOperationSetups--;
_checkIdleAndDisconnect();
}
return completer.future;
}
Expand Down Expand Up @@ -873,6 +906,7 @@ class WebSocketTransport implements DataConnectTransport {
_sendPendingSubscriptions(authToken, appCheckToken);
} finally {
_pendingOperationSetups--;
_checkIdleAndDisconnect();
}
},
onCancel: () {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ import 'package:mockito/annotations.dart';
import 'package:mockito/mockito.dart';

@GenerateNiceMocks([MockSpec<FirebaseApp>(), MockSpec<ConnectorConfig>()])
import '../firebase_data_connect_test.mocks.dart';
import 'cache_manager_test.mocks.dart';
import '../network/rest_transport_test.mocks.dart';

class MockTransportOptions extends Mock implements TransportOptions {}
Expand Down
Original file line number Diff line number Diff line change
@@ -1,17 +1,3 @@
// Copyright 2026 Google LLC
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

// Mocks generated by Mockito 5.4.6 from annotations
// in firebase_data_connect/test/src/cache/cache_manager_test.dart.
// Do not manually edit this file.
Expand Down Expand Up @@ -123,7 +109,7 @@ class MockFirebaseApp extends _i1.Mock implements _i3.FirebaseApp {

@override
void registerService<T extends _i3.FirebaseService>(
T service, {
T? service, {
_i5.Future<void> Function(T)? dispose,
}) =>
super.noSuchMethod(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ import 'package:flutter_test/flutter_test.dart';
import 'package:mockito/annotations.dart';
import 'package:mockito/mockito.dart';

@GenerateNiceMocks([MockSpec<FirebaseApp>(), MockSpec<ConnectorConfig>()])
@GenerateNiceMocks([MockSpec<ConnectorConfig>()])
import 'firebase_data_connect_test.mocks.dart';

class MockFirebaseAuth extends Mock implements FirebaseAuth {
Expand Down Expand Up @@ -75,24 +75,25 @@ class MockQueryManager extends Mock implements QueryManager {}

void main() {
group('FirebaseDataConnect', () {
late MockFirebaseApp mockApp;
late DynamicMockFirebaseApp mockApp;
late MockFirebaseAuth mockAuth;
late MockFirebaseAppCheck mockAppCheck;
late MockConnectorConfig mockConnectorConfig;

setUp(() {
mockApp = MockFirebaseApp();
mockAuth = MockFirebaseAuth();
mockAppCheck = MockFirebaseAppCheck();
mockConnectorConfig = MockConnectorConfig();

when(mockApp.options).thenReturn(
const FirebaseOptions(
mockApp = DynamicMockFirebaseApp(
name: 'default',
options: const FirebaseOptions(
apiKey: 'fake_api_key',
appId: 'fake_app_id',
messagingSenderId: 'fake_messaging_sender_id',
projectId: 'fake_project_id',
),
mockAuth: mockAuth,
mockAppCheck: mockAppCheck,
);
when(mockConnectorConfig.location).thenReturn('us-central1');
when(mockConnectorConfig.connector).thenReturn('connector');
Expand Down Expand Up @@ -181,7 +182,12 @@ void main() {
test('instanceFor returns cached instance if available', () {
FirebaseDataConnect.cachedInstances.clear(); // Clear cache first

when(mockApp.name).thenReturn('appName');
mockApp = DynamicMockFirebaseApp(
name: 'appName',
options: mockApp.options,
mockAuth: mockAuth,
mockAppCheck: mockAppCheck,
);
when(mockConnectorConfig.toJson()).thenReturn('connectorConfigStr');

final dataConnect = FirebaseDataConnect(
Expand All @@ -208,7 +214,12 @@ void main() {
test('instanceFor creates new instance if not cached', () {
FirebaseDataConnect.cachedInstances.clear(); // Clear cache first

when(mockApp.name).thenReturn('appName');
mockApp = DynamicMockFirebaseApp(
name: 'appName',
options: mockApp.options,
mockAuth: mockAuth,
mockAppCheck: mockAppCheck,
);
when(mockConnectorConfig.toJson()).thenReturn('connectorConfigStr');

final instance = FirebaseDataConnect.instanceFor(
Expand Down
Loading
Loading