diff --git a/packages/realtime_client/lib/src/realtime_client.dart b/packages/realtime_client/lib/src/realtime_client.dart index ba6ed51ef..b1986f931 100644 --- a/packages/realtime_client/lib/src/realtime_client.dart +++ b/packages/realtime_client/lib/src/realtime_client.dart @@ -138,7 +138,7 @@ class RealtimeClient { final RealtimeEncode encode; final RealtimeDecode decode; late TimerCalculation reconnectAfterMs; - WebSocketChannel? conn; + WebSocketChannel? connection; StreamSubscription? _connectionSubscription; List sendBuffer = []; Map> stateChangeCallbacks = { @@ -153,7 +153,7 @@ class RealtimeClient { @Deprecated("No longer used. Will be removed in the next major version.") int longpollerTimeout = 20000; - SocketStates? connState; + SocketStates? connectionStatus; Future Function()? customAccessToken; /// Initializes the Socket @@ -256,8 +256,8 @@ class RealtimeClient { /// Connects the socket. @internal Future connect() async { - if (conn != null) { - if (connState != SocketStates.closed) { + if (connection != null) { + if (connectionStatus != SocketStates.closed) { return; } await disconnect(); @@ -266,52 +266,52 @@ class RealtimeClient { try { log('transport', 'connecting to $endPointURL', null); log('transport', 'connecting', null, Level.FINE); - connState = SocketStates.connecting; - final WebSocketChannel localConn = transport(endPointURL, headers); - conn = localConn; + connectionStatus = SocketStates.connecting; + final WebSocketChannel localConnection = transport(endPointURL, headers); + connection = localConnection; try { - await localConn.ready; + await localConnection.ready; } catch (error) { // Bail out if disconnect() ran or a new connect() started during await - if (conn != localConn) { + if (connection != localConnection) { return; } // Don't schedule a reconnect and emit error if connection has been // closed by the user or [disconnect] waits for the connection to be // ready before closing it. - if (connState != SocketStates.disconnected && - connState != SocketStates.disconnecting) { - connState = SocketStates.closed; - _onConnError(error); + if (connectionStatus != SocketStates.disconnected && + connectionStatus != SocketStates.disconnecting) { + connectionStatus = SocketStates.closed; + _onConnectionError(error); reconnectTimer.scheduleTimeout(); } return; } // Guard: bail out if disconnect() ran during the await - if (conn != localConn || connState != SocketStates.connecting) { + if (connection != localConnection || connectionStatus != SocketStates.connecting) { return; } - connState = SocketStates.open; + connectionStatus = SocketStates.open; - _onConnOpen(); - _connectionSubscription = localConn.stream.listen( - (message) => onConnMessage(message), - onError: _onConnError, + _onConnectionOpen(); + _connectionSubscription = localConnection.stream.listen( + (message) => onConnectionMessage(message), + onError: _onConnectionError, onDone: () { // communication has been closed - if (connState != SocketStates.disconnected && - connState != SocketStates.disconnecting) { - connState = SocketStates.closed; + if (connectionStatus != SocketStates.disconnected && + connectionStatus != SocketStates.disconnecting) { + connectionStatus = SocketStates.closed; } - _onConnClose(); + _onConnectionClose(); }, ); } catch (e) { /// General error handling - _onConnError(e); + _onConnectionError(e); } } @@ -323,14 +323,14 @@ class RealtimeClient { /// Disconnects the socket with status [code] and [reason] for the disconnect Future disconnect({int? code, String? reason}) async { _cancelPendingDisconnect(); - final conn = this.conn; - if (conn != null) { - final oldState = connState; + final connection = this.connection; + if (connection != null) { + final oldState = connectionStatus; final shouldCloseSink = oldState == SocketStates.open || oldState == SocketStates.connecting; if (shouldCloseSink) { // Don't set the state to `disconnecting` if the connection is already closed. - connState = SocketStates.disconnecting; + connectionStatus = SocketStates.disconnecting; log('transport', 'disconnecting', { 'code': code, 'reason': reason, @@ -349,34 +349,34 @@ class RealtimeClient { // avoid hanging the client. This is done by mimicking the onDone // callback of the connection stream. By canceling the subscription, // we avoid calling the onDone too. - connState = SocketStates.disconnected; - _onConnClose(); + connectionStatus = SocketStates.disconnected; + _onConnectionClose(); } if (code != null) { // Add a timeout to close the sink to avoid hanging in case something // is wrong with the connection. // The Dart SDK has a timeout of 5 seconds for closing the IO WebSocket connection, so we set a timeout of 6 seconds here to avoid hanging indefinitely. - await conn.sink + await connection.sink .close(code, reason ?? '') .timeout(connectionCloseTimeout, onTimeout: onTimeout); } else { - await conn.sink.close().timeout( + await connection.sink.close().timeout( connectionCloseTimeout, onTimeout: onTimeout, ); } - connState = SocketStates.disconnected; + connectionStatus = SocketStates.disconnected; log('transport', 'disconnected', null, Level.FINE); } - // Cancel any reconnect scheduled by `_onConnClose`. When the socket has - // already dropped (`connState == closed`) the block above is skipped, so + // Cancel any reconnect scheduled by `_onConnectionClose`. When the socket has + // already dropped (`connectionStatus == closed`) the block above is skipped, so // without this an armed backoff timer would fire after the user // explicitly disconnected and silently reopen the connection. reconnectTimer.cancel(); - this.conn = null; + this.connection = null; await _connectionSubscription?.cancel(); _connectionSubscription = null; @@ -444,7 +444,7 @@ class RealtimeClient { _heartbeatController.stream; /// Returns the current state of the socket. - String get connectionState => switch (connState) { + String get connectionState => switch (connectionStatus) { SocketStates.connecting => 'connecting', SocketStates.open => 'open', SocketStates.disconnecting => 'disconnecting', @@ -453,7 +453,7 @@ class RealtimeClient { }; /// Returns `true` is the connection is open. - bool get isConnected => connState == SocketStates.open; + bool get isConnected => connectionStatus == SocketStates.open; /// Removes a subscription from the socket. @internal @@ -522,7 +522,7 @@ class RealtimeClient { // ignore: function-always-returns-null String? push(Message message) { void callback() { - conn?.sink.add(encode(message.toJson())); + connection?.sink.add(encode(message.toJson())); } log( @@ -539,7 +539,7 @@ class RealtimeClient { return null; } - void onConnMessage(Object rawMessage) { + void onConnectionMessage(Object rawMessage) { final Map message; try { message = decode(rawMessage); @@ -645,7 +645,7 @@ class RealtimeClient { } } - void _onConnOpen() { + void _onConnectionOpen() { log('transport', 'connected to $endPointURL'); log('transport', 'connected', null, Level.FINE); unawaited(_resolveAccessTokenAndFlush()); @@ -672,17 +672,17 @@ class RealtimeClient { } /// communication has been closed - void _onConnClose() { - final statusCode = conn?.closeCode; + void _onConnectionClose() { + final statusCode = connection?.closeCode; RealtimeCloseEvent? event; if (statusCode != null) { - event = RealtimeCloseEvent(code: statusCode, reason: conn?.closeReason); + event = RealtimeCloseEvent(code: statusCode, reason: connection?.closeReason); } log('transport', 'close', event, Level.FINE); /// SocketStates.disconnected: by user with socket.disconnect() /// SocketStates.closed: NOT by user, should try to reconnect - if (connState == SocketStates.closed) { + if (connectionStatus == SocketStates.closed) { _triggerChanError(event); reconnectTimer.scheduleTimeout(); } @@ -692,7 +692,7 @@ class RealtimeClient { } } - void _onConnError(dynamic error) { + void _onConnectionError(dynamic error) { log('transport', error.toString()); _triggerChanError(error); for (final callback in stateChangeCallbacks['error']!) { @@ -777,7 +777,7 @@ class RealtimeClient { 'heartbeat timeout. Attempting to re-establish conn', ); _heartbeatController.add(RealtimeHeartbeatStatus.timeout); - unawaited(conn?.sink.close(Constants.wsCloseNormal, 'heartbeat timeout')); + unawaited(connection?.sink.close(Constants.wsCloseNormal, 'heartbeat timeout')); return; } pendingHeartbeatRef = makeRef(); diff --git a/packages/realtime_client/test/heartbeat_test.dart b/packages/realtime_client/test/heartbeat_test.dart index 979b54f64..1233fb670 100644 --- a/packages/realtime_client/test/heartbeat_test.dart +++ b/packages/realtime_client/test/heartbeat_test.dart @@ -40,7 +40,7 @@ void main() { }); test('emits sent when a heartbeat is pushed', () async { - client.connState = SocketStates.open; + client.connectionStatus = SocketStates.open; await client.sendHeartbeat(); await pumpEventQueue(); @@ -52,7 +52,7 @@ void main() { test( 'emits timeout when the previous heartbeat was not acknowledged', () async { - client.connState = SocketStates.open; + client.connectionStatus = SocketStates.open; client.pendingHeartbeatRef = 'stale-ref'; await client.sendHeartbeat(); @@ -66,7 +66,7 @@ void main() { test('emits ok when the heartbeat reply succeeds', () async { client.pendingHeartbeatRef = 'ref-1'; - client.onConnMessage(heartbeatReply('ref-1', 'ok')); + client.onConnectionMessage(heartbeatReply('ref-1', 'ok')); await pumpEventQueue(); expect(statuses, [RealtimeHeartbeatStatus.ok]); @@ -76,7 +76,7 @@ void main() { test('emits error when the heartbeat reply fails', () async { client.pendingHeartbeatRef = 'ref-2'; - client.onConnMessage(heartbeatReply('ref-2', 'error')); + client.onConnectionMessage(heartbeatReply('ref-2', 'error')); await pumpEventQueue(); expect(statuses, [RealtimeHeartbeatStatus.error]); @@ -87,7 +87,7 @@ void main() { () async { client.pendingHeartbeatRef = 'ref-3'; - client.onConnMessage(heartbeatReply('other-ref', 'ok')); + client.onConnectionMessage(heartbeatReply('other-ref', 'ok')); await pumpEventQueue(); expect(statuses, isEmpty); diff --git a/packages/realtime_client/test/socket_test.dart b/packages/realtime_client/test/socket_test.dart index 1c58dfbe4..b7d3206b6 100644 --- a/packages/realtime_client/test/socket_test.dart +++ b/packages/realtime_client/test/socket_test.dart @@ -200,12 +200,12 @@ void main() { test('establishes websocket connection with endpoint', () async { final connectFuture = socket.connect(); - expect(socket.connState, SocketStates.connecting); + expect(socket.connectionStatus, SocketStates.connecting); - final connection = socket.conn; + final connection = socket.connection; await connectFuture; - expect(socket.connState, SocketStates.open); + expect(socket.connectionStatus, SocketStates.open); expect(connection, isA()); //! Not verifying connection url @@ -253,9 +253,9 @@ void main() { test('is idempotent', () { unawaited(socket.connect()); - final connection = socket.conn; + final connection = socket.connection; unawaited(socket.connect()); - expect(socket.conn, connection); + expect(socket.connection, connection); }); }); @@ -272,10 +272,10 @@ void main() { test('removes existing connection', () async { await socket.connect(); - expect(socket.conn, isNotNull); + expect(socket.connection, isNotNull); await socket.disconnect(); - expect(socket.conn, isNull); + expect(socket.connection, isNull); }); test('calls connection close callback', () async { @@ -297,7 +297,7 @@ void main() { const tReason = 'reason'; await mockedSocket.connect(); - mockedSocket.connState = SocketStates.open; + mockedSocket.connectionStatus = SocketStates.open; await Future.delayed(const Duration(milliseconds: 200)); await mockedSocket.disconnect(code: tCode, reason: tReason); await Future.delayed(const Duration(milliseconds: 200)); @@ -312,19 +312,19 @@ void main() { test('disconnecting a closed connections stays closed', () async { await socket.connect(); - expect(socket.connState, SocketStates.open); + expect(socket.connectionStatus, SocketStates.open); await mockServer.close(); await Future.delayed(const Duration(milliseconds: 200)); - expect(socket.connState, SocketStates.closed); - expect(socket.conn, isNotNull); + expect(socket.connectionStatus, SocketStates.closed); + expect(socket.connection, isNotNull); final disconnectFuture = socket.disconnect(); - // `connState` stays `closed` during disconnect - expect(socket.connState, SocketStates.closed); + // `connectionStatus` stays `closed` during disconnect + expect(socket.connectionStatus, SocketStates.closed); await disconnectFuture; - expect(socket.connState, SocketStates.closed); - expect(socket.conn, isNull); + expect(socket.connectionStatus, SocketStates.closed); + expect(socket.connection, isNull); }); test('cancels a pending reconnect after an unexpected drop', () async { @@ -361,7 +361,7 @@ void main() { // is marked closed and a reconnect is scheduled. await streamController.close(); await Future.delayed(const Duration(milliseconds: 5)); - expect(mockedSocket.connState, SocketStates.closed); + expect(mockedSocket.connectionStatus, SocketStates.closed); // The user disconnects explicitly while the socket is already closed. await mockedSocket.disconnect(); @@ -415,22 +415,22 @@ void main() { await mockedSocket.connect(); expect(connectCount, 1); - expect(mockedSocket.connState, SocketStates.open); + expect(mockedSocket.connectionStatus, SocketStates.open); // Simulate the server dropping the connection. await firstController.close(); await Future.delayed(const Duration(milliseconds: 5)); - expect(mockedSocket.connState, SocketStates.closed); + expect(mockedSocket.connectionStatus, SocketStates.closed); // A manual reconnect must open a fresh connection instead of being a - // no-op because `conn` still references the dropped socket. + // no-op because `connection` still references the dropped socket. await mockedSocket.connect(); expect( connectCount, 2, reason: 'manual connect() must reconnect after a drop', ); - expect(mockedSocket.connState, SocketStates.open); + expect(mockedSocket.connectionStatus, SocketStates.open); await mockedSocket.disconnect(); }); @@ -486,22 +486,22 @@ void main() { // The retry counter must grow (1, 2, 3, ...) across reconnect attempts // instead of being reset to 1 on every `disconnect()` in `_reconnect`. expect(triesSeen.take(3), [1, 2, 3]); - expect(mockedSocket.connState, SocketStates.open); + expect(mockedSocket.connectionStatus, SocketStates.open); await mockedSocket.disconnect(); }); test('disconnecting an open connection', () async { await socket.connect(); - expect(socket.connState, SocketStates.open); + expect(socket.connectionStatus, SocketStates.open); final disconnectFuture = socket.disconnect(); - // `connState` stays `closed` during disconnect - expect(socket.connState, SocketStates.disconnecting); + // `connectionStatus` stays `closed` during disconnect + expect(socket.connectionStatus, SocketStates.disconnecting); await disconnectFuture; - expect(socket.connState, SocketStates.disconnected); - expect(socket.conn, isNull); + expect(socket.connectionStatus, SocketStates.disconnected); + expect(socket.connection, isNull); }); test('does not throw when no connection', () { @@ -528,11 +528,11 @@ void main() { when(() => mockedSink.close()).thenAnswer((_) => closeCompleter.future); await mockedSocket.connect(); - expect(mockedSocket.connState, SocketStates.open); + expect(mockedSocket.connectionStatus, SocketStates.open); await mockedSocket.disconnect(); - expect(mockedSocket.connState, SocketStates.disconnected); - expect(mockedSocket.conn, isNull); + expect(mockedSocket.connectionStatus, SocketStates.disconnected); + expect(mockedSocket.connection, isNull); expect(closeCallbacks, 1); verify(() => mockedSink.close()).called(1); @@ -813,7 +813,7 @@ void main() { test('sends data to connection when connected', () { unawaited(mockedSocket.connect()); - mockedSocket.connState = SocketStates.open; + mockedSocket.connectionStatus = SocketStates.open; final message = Message( topic: topic, @@ -830,7 +830,7 @@ void main() { test('buffers data when not connected', () async { unawaited(mockedSocket.connect()); - mockedSocket.connState = SocketStates.connecting; + mockedSocket.connectionStatus = SocketStates.connecting; expect(mockedSocket.sendBuffer, isEmpty); @@ -854,7 +854,7 @@ void main() { test('sends a broadcast with a binary payload as a binary frame', () { unawaited(mockedSocket.connect()); - mockedSocket.connState = SocketStates.open; + mockedSocket.connectionStatus = SocketStates.open; final binaryPayload = Uint8List.fromList([1, 2, 3]); final message = Message( @@ -886,7 +886,7 @@ void main() { version: RealtimeProtocolVersion.v1, ); unawaited(legacySocket.connect()); - legacySocket.connState = SocketStates.open; + legacySocket.connectionStatus = SocketStates.open; final legacyData = json.encode({ 'topic': topic, @@ -921,7 +921,7 @@ void main() { encode: (_) => 'custom-frame', ); unawaited(customSocket.connect()); - customSocket.connState = SocketStates.open; + customSocket.connectionStatus = SocketStates.open; customSocket.push( Message(topic: topic, payload: payload, event: event, ref: ref), @@ -933,11 +933,11 @@ void main() { }); }); - group('onConnMessage', () { + group('onConnectionMessage', () { test('drops a malformed frame without throwing', () { final socket = RealtimeClient(socketEndpoint); expect( - () => socket.onConnMessage('{"not": "an array"}'), + () => socket.onConnectionMessage('{"not": "an array"}'), returnsNormally, ); }); @@ -966,7 +966,7 @@ void main() { ...payload, ]); - socket.onConnMessage(frame); + socket.onConnectionMessage(frame); expect(received, { 'type': 'broadcast', @@ -990,7 +990,7 @@ void main() { callback: (payload) => received = payload, ); - socket.onConnMessage( + socket.onConnectionMessage( json.encode({ 'topic': 'realtime:room', 'event': 'broadcast', @@ -1193,7 +1193,7 @@ void main() { verify(() => erroredChannel.rejoin()).called(1); verifyNever(() => healthyChannel.rejoin()); expect(opens, 1); - expect(socket.connState, SocketStates.open); + expect(socket.connectionStatus, SocketStates.open); await socket.disconnect(); await streamController.close(); @@ -1295,7 +1295,7 @@ void main() { //! Unimplemented Test: closes socket when heartbeat is not ack'd within heartbeat window test('pushes heartbeat data when connected', () async { - mockedSocket.connState = SocketStates.open; + mockedSocket.connectionStatus = SocketStates.open; await mockedSocket.sendHeartbeat(); @@ -1303,7 +1303,7 @@ void main() { }); test('no ops when not connected', () async { - mockedSocket.connState = SocketStates.connecting; + mockedSocket.connectionStatus = SocketStates.connecting; await mockedSocket.sendHeartbeat(); verifyNever(() => mockedSink.add(any())); @@ -1344,13 +1344,13 @@ void main() { await connectFuture; // Should NOT have transitioned to open because disconnect nullified connection - expect(socket.connState, isNot(SocketStates.open)); - expect(socket.conn, isNull); + expect(socket.connectionStatus, isNot(SocketStates.open)); + expect(socket.connection, isNull); }, ); test( - 'connect bails out when connState changes during await ready', + 'connect bails out when connectionStatus changes during await ready', () async { final readyCompleter = Completer(); final mockedSocketChannel = MockIOWebSocketChannel(); @@ -1381,7 +1381,7 @@ void main() { await disconnectFuture; await connectFuture; - expect(socket.connState, isNot(SocketStates.open)); + expect(socket.connectionStatus, isNot(SocketStates.open)); }, ); @@ -1441,8 +1441,8 @@ void main() { readyCompleter2.complete(); await socket.connect(); - expect(socket.connState, SocketStates.open); - expect(socket.conn, mockedSocketChannel2); + expect(socket.connectionStatus, SocketStates.open); + expect(socket.connection, mockedSocketChannel2); await socket.disconnect(); await streamController2.close(); diff --git a/packages/supabase_flutter/lib/src/supabase.dart b/packages/supabase_flutter/lib/src/supabase.dart index 0f4eeb091..a029548cb 100644 --- a/packages/supabase_flutter/lib/src/supabase.dart +++ b/packages/supabase_flutter/lib/src/supabase.dart @@ -312,7 +312,7 @@ class Supabase { // paused or detached — disconnect the WebSocket if it is active. // These states are not triggered on web if (realtime.isConnected || - realtime.connState == SocketStates.connecting) { + realtime.connectionStatus == SocketStates.connecting) { await realtime.disconnect(); } } diff --git a/packages/supabase_flutter/test/lifecycle_test.dart b/packages/supabase_flutter/test/lifecycle_test.dart index 923b132ac..6047b12a2 100644 --- a/packages/supabase_flutter/test/lifecycle_test.dart +++ b/packages/supabase_flutter/test/lifecycle_test.dart @@ -138,7 +138,7 @@ void main() { // Connect with ready completed immediately await connectAndReady(realtime); - expect(realtime.connState, SocketStates.open); + expect(realtime.connectionStatus, SocketStates.open); // paused → triggers disconnect binding.handleAppLifecycleStateChanged(AppLifecycleState.inactive); @@ -152,8 +152,8 @@ void main() { // Complete any pending ready futures (reconnect) await settleLifecycle(); - expect(realtime.connState, SocketStates.open); - expect(realtime.conn, isNotNull); + expect(realtime.connectionStatus, SocketStates.open); + expect(realtime.connection, isNotNull); }); test('paused → resumed → inactive → resumed ' @@ -164,7 +164,7 @@ void main() { realtime.channel('test'); await connectAndReady(realtime); - expect(realtime.connState, SocketStates.open); + expect(realtime.connectionStatus, SocketStates.open); // paused → starts disconnect binding.handleAppLifecycleStateChanged(AppLifecycleState.inactive); @@ -185,8 +185,8 @@ void main() { await settleLifecycle(); // Should have reconnected, not stuck disconnecting - expect(realtime.connState, SocketStates.open); - expect(realtime.conn, isNotNull); + expect(realtime.connectionStatus, SocketStates.open); + expect(realtime.connection, isNotNull); }); test('rapid paused → resumed → paused → resumed ' @@ -197,7 +197,7 @@ void main() { realtime.channel('test'); await connectAndReady(realtime); - expect(realtime.connState, SocketStates.open); + expect(realtime.connectionStatus, SocketStates.open); // Rapid lifecycle flapping binding.handleAppLifecycleStateChanged(AppLifecycleState.inactive); @@ -217,8 +217,8 @@ void main() { // Complete all pending ready futures as they appear await settleLifecycle(); - expect(realtime.connState, SocketStates.open); - expect(realtime.conn, isNotNull); + expect(realtime.connectionStatus, SocketStates.open); + expect(realtime.connection, isNotNull); }); test('resumed then paused before connect completes ' @@ -229,7 +229,7 @@ void main() { realtime.channel('test'); await connectAndReady(realtime); - expect(realtime.connState, SocketStates.open); + expect(realtime.connectionStatus, SocketStates.open); // paused → triggers disconnect binding.handleAppLifecycleStateChanged(AppLifecycleState.inactive); @@ -250,8 +250,8 @@ void main() { await settleLifecycle(); // Should be disconnected since the last event was paused - expect(realtime.connState, SocketStates.disconnected); - expect(realtime.conn, isNull); + expect(realtime.connectionStatus, SocketStates.disconnected); + expect(realtime.connection, isNull); }); }); }