Skip to content
Merged
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
2 changes: 1 addition & 1 deletion packages/realtime_client/lib/src/push.dart
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,7 @@ class Push {
});
}

void trigger(String status, dynamic response) {
void trigger(String status, Map<String, dynamic> response) {
if (_refEvent != null) {
_channel.trigger(_refEvent!, {'status': status, 'response': response});
}
Expand Down
142 changes: 75 additions & 67 deletions packages/realtime_client/lib/src/realtime_channel.dart
Original file line number Diff line number Diff line change
Expand Up @@ -171,74 +171,13 @@ class RealtimeChannel {
joinedOnce = true;
rejoin(timeout ?? _timeout);

joinPush.receive(
joinPush
.receive(
'ok',
(response) async {
final serverPostgresFilters = response['postgres_changes'];
if (socket.accessToken != null) {
try {
// ignore: avoid-passing-self-as-argument
await socket.setAuth(socket.accessToken);
} on FormatException catch (e) {
// The cached access token may have expired by the time the
// channel rejoins (e.g. after the device wakes from a long
// sleep). Auth state listeners will re-call setAuth with a
// fresh token shortly after, so swallow this specific error
// to avoid surfacing it as an uncaught exception. The same
// filter is applied in `SupabaseClient._handleTokenChanged`.
if (!e.message.contains('InvalidJWTToken')) {
rethrow;
}
}
}

if (serverPostgresFilters == null) {
if (callback != null) {
callback(RealtimeSubscribeStatus.subscribed, null);
}
return;
}
final clientPostgresBindings = _bindings['postgres_changes'];
final bindingsLen = clientPostgresBindings?.length ?? 0;
final newPostgresBindings = <Binding>[];

for (var i = 0; i < bindingsLen; i++) {
final clientPostgresBinding = clientPostgresBindings![i];

final event = clientPostgresBinding.filter['event'];
final schema = clientPostgresBinding.filter['schema'];
final table = clientPostgresBinding.filter['table'];
final filter = clientPostgresBinding.filter['filter'];
final serverPostgresFilter = serverPostgresFilters[i];

if (serverPostgresFilter != null &&
serverPostgresFilter['event'] == event &&
serverPostgresFilter['schema'] == schema &&
serverPostgresFilter['table'] == table &&
serverPostgresFilter['filter'] == filter) {
newPostgresBindings.add(clientPostgresBinding.copyWith(
id: serverPostgresFilter['id']?.toString(),
));
} else {
unawaited(unsubscribe());
if (callback != null) {
callback(
RealtimeSubscribeStatus.channelError,
Exception(
'mismatch between server and client bindings for postgres changes'),
);
}
return;
}
}

_bindings['postgres_changes'] = newPostgresBindings;

if (callback != null) {
callback(RealtimeSubscribeStatus.subscribed, null);
}
},
).receive('error', (error) {
(response) =>
unawaited(_handleJoinOk(response as Map<String, dynamic>, callback)),
)
.receive('error', (error) {
if (callback != null) {
callback(
RealtimeSubscribeStatus.channelError,
Expand All @@ -255,6 +194,75 @@ class RealtimeChannel {
return this;
}

Future<void> _handleJoinOk(
Map<String, dynamic> response,
void Function(RealtimeSubscribeStatus status, Object? error)? callback,
) async {
final serverPostgresFilters = response['postgres_changes'];
if (socket.accessToken != null) {
try {
// ignore: avoid-passing-self-as-argument
await socket.setAuth(socket.accessToken);
} on FormatException catch (e) {
// The cached access token may have expired by the time the
// channel rejoins (e.g. after the device wakes from a long
// sleep). Auth state listeners will re-call setAuth with a
// fresh token shortly after, so swallow this specific error
// to avoid surfacing it as an uncaught exception. The same
// filter is applied in `SupabaseClient._handleTokenChanged`.
if (!e.message.contains('InvalidJWTToken')) {
rethrow;
}
}
}

if (serverPostgresFilters == null) {
if (callback != null) {
callback(RealtimeSubscribeStatus.subscribed, null);
}
return;
}
final clientPostgresBindings = _bindings['postgres_changes'];
final bindingsLen = clientPostgresBindings?.length ?? 0;
final newPostgresBindings = <Binding>[];

for (var i = 0; i < bindingsLen; i++) {
final clientPostgresBinding = clientPostgresBindings![i];

final event = clientPostgresBinding.filter['event'];
final schema = clientPostgresBinding.filter['schema'];
final table = clientPostgresBinding.filter['table'];
final filter = clientPostgresBinding.filter['filter'];
final serverPostgresFilter = serverPostgresFilters[i];

if (serverPostgresFilter != null &&
serverPostgresFilter['event'] == event &&
serverPostgresFilter['schema'] == schema &&
serverPostgresFilter['table'] == table &&
serverPostgresFilter['filter'] == filter) {
newPostgresBindings.add(clientPostgresBinding.copyWith(
id: serverPostgresFilter['id']?.toString(),
));
} else {
unawaited(unsubscribe());
if (callback != null) {
callback(
RealtimeSubscribeStatus.channelError,
Exception(
'mismatch between server and client bindings for postgres changes'),
);
}
return;
}
}

_bindings['postgres_changes'] = newPostgresBindings;

if (callback != null) {
callback(RealtimeSubscribeStatus.subscribed, null);
}
}

List<SinglePresenceState> presenceState() {
return presence.state.entries
.map((entry) =>
Expand Down
12 changes: 7 additions & 5 deletions packages/realtime_client/lib/src/realtime_client.dart
Original file line number Diff line number Diff line change
Expand Up @@ -204,10 +204,7 @@ class RealtimeClient {
this.reconnectAfterMs =
reconnectAfterMs ?? RetryTimer.createRetryFunction();
reconnectTimer = RetryTimer(
() async {
await disconnect();
await connect();
},
() => unawaited(_reconnect()),
this.reconnectAfterMs,
);
}
Expand Down Expand Up @@ -271,6 +268,11 @@ class RealtimeClient {
}
}

Future<void> _reconnect() async {
await disconnect();
await connect();
}

/// Disconnects the socket with status [code] and [reason] for the disconnect
Future<void> disconnect({int? code, String? reason}) async {
final conn = this.conn;
Expand Down Expand Up @@ -526,7 +528,7 @@ class RealtimeClient {
if (heartbeatTimer != null) heartbeatTimer!.cancel();
heartbeatTimer = Timer.periodic(
Duration(milliseconds: heartbeatIntervalMs),
(Timer t) async => await sendHeartbeat(),
(Timer t) => unawaited(sendHeartbeat()),
);
for (final callback in stateChangeCallbacks['open']!) {
callback();
Expand Down
99 changes: 52 additions & 47 deletions packages/realtime_client/test/channel_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -392,66 +392,71 @@ void main() {
});

test('send message via ws conn when subscribed to channel', () async {
channel.subscribe((status, [error]) async {
final subscribed = Completer<void>();
channel.subscribe((status, [error]) {
if (status == RealtimeSubscribeStatus.subscribed) {
final completer = Completer<ChannelResponse>();
unawaited(
channel.send(
type: RealtimeListenTypes.broadcast,
payload: {
'myKey': 'myValue',
},
).then(
(value) => completer.complete(value),
onError: (Object e, StackTrace stackTrace) =>
completer.completeError(e, stackTrace),
),
);

await for (final HttpRequest req in mockServer) {
expect(req.uri.toString(), startsWith('/realtime/v1/websocket'));
await req.response.close();
break;
}
expect(await completer.future, ChannelResponse.ok);
subscribed.complete();
}
});

// Accept the websocket the client opens on subscribe, then reply to the
// channel join so it transitions to subscribed.
final serverSocket =
await mockServer.first.then(WebSocketTransformer.upgrade);
final broadcast = Completer<List<dynamic>>();
serverSocket.listen((frame) {
final message = jsonDecode(frame as String) as List;
switch (message[3]) {
case 'phx_join':
serverSocket.add(jsonEncode([
message[0],
message[1],
message[2],
'phx_reply',
{'status': 'ok', 'response': <String, dynamic>{}},
]));
case 'broadcast':
broadcast.complete(message);
}
});
await subscribed.future;

// Once subscribed, broadcasts are pushed over the websocket instead of
// falling back to the REST endpoint.
final sendResult = await channel.send(
type: RealtimeListenTypes.broadcast,
payload: {'myKey': 'myValue'},
);
expect(sendResult, ChannelResponse.ok);

final message = await broadcast.future;
expect(message[2], 'realtime:myTopic');
expect(message[3], 'broadcast');
expect(message[4], containsPair('myKey', 'myValue'));
});

test(
'send message via http request to Broadcast endpoint when not subscribed to channel',
() async {
final completer = Completer<ChannelResponse>();
unawaited(
channel.send(
type: RealtimeListenTypes.broadcast,
payload: {
'myKey': 'myValue',
},
).then(
(value) => completer.complete(value),
onError: (Object e, StackTrace stackTrace) =>
completer.completeError(e, stackTrace),
),
final requestFuture = mockServer.first;
final sendFuture = channel.send(
type: RealtimeListenTypes.broadcast,
payload: {'myKey': 'myValue'},
);

await for (final HttpRequest req in mockServer) {
expect(req.uri.toString(), '/realtime/v1/api/broadcast');
expect(req.headers.value('apikey'), 'supabaseKey');
final request = await requestFuture;
expect(request.uri.toString(), '/realtime/v1/api/broadcast');
expect(request.headers.value('apikey'), 'supabaseKey');

final body = json.decode(await utf8.decodeStream(req));
final message = body['messages'].first;
final payload = message['payload'];
final private = message['private'];
final body = json.decode(await utf8.decodeStream(request));
final message = body['messages'].first;
expect(message['payload'], containsPair('myKey', 'myValue'));
expect(message, containsPair('topic', 'myTopic'));
expect(message['private'], isTrue);

expect(payload, containsPair('myKey', 'myValue'));
expect(message, containsPair('topic', 'myTopic'));
expect(private, isTrue);
await request.response.close();

await req.response.close();
break;
}
expect(await completer.future, ChannelResponse.ok);
expect(await sendFuture, ChannelResponse.ok);
});
});

Expand Down
Loading
Loading