From 623d48e32bfb3d08fd70168fcb1e16f695a98033 Mon Sep 17 00:00:00 2001 From: Sahil Kumar Date: Fri, 26 Feb 2021 18:17:43 +0530 Subject: [PATCH] Add custom sort support for offline channels Signed-off-by: Sahil Kumar --- packages/stream_chat/lib/src/api/channel.dart | 4 +- .../stream_chat/lib/src/api/requests.dart | 14 +- packages/stream_chat/lib/src/client.dart | 253 ++++++++---------- .../lib/src/db/chat_persistence_client.dart | 30 +-- .../lib/src/models/channel_state.dart | 2 +- packages/stream_chat/test/version_test.dart | 120 ++++++++- .../lib/src/channel_info.dart | 5 +- .../lib/src/channel_list_view.dart | 168 +++--------- .../lib/src/channel_list_core.dart | 79 ++---- .../lib/src/channels_bloc.dart | 22 +- .../example/pubspec.yaml | 2 +- .../lib/src/dao/channel_query_dao.dart | 130 ++++----- .../src/stream_chat_persistence_client.dart | 7 +- 13 files changed, 394 insertions(+), 442 deletions(-) diff --git a/packages/stream_chat/lib/src/api/channel.dart b/packages/stream_chat/lib/src/api/channel.dart index a7231235..1f5da466 100644 --- a/packages/stream_chat/lib/src/api/channel.dart +++ b/packages/stream_chat/lib/src/api/channel.dart @@ -1245,7 +1245,9 @@ class ChannelClientState { _channel._client.chatPersistenceClient ?.getChannelStateByCid(_channel.cid) ?.then((state) { - updateChannelState(state); + // Replacing the persistence state members with the latest `channelState.members` + // as they may have changes over the time. + updateChannelState(state.copyWith(members: channelState.members)); retryFailedMessages(); }); }); diff --git a/packages/stream_chat/lib/src/api/requests.dart b/packages/stream_chat/lib/src/api/requests.dart index fdf1f34f..414113de 100644 --- a/packages/stream_chat/lib/src/api/requests.dart +++ b/packages/stream_chat/lib/src/api/requests.dart @@ -4,7 +4,7 @@ part 'requests.g.dart'; /// Sorting options @JsonSerializable(createFactory: false) -class SortOption { +class SortOption { /// Ascending order static const ASC = 1; @@ -17,6 +17,10 @@ class SortOption { /// A sorting direction final int direction; + /// Sorting field Comparator required for offline sorting + @JsonKey(ignore: true) + final Comparator comparator; + /// Creates a new SortOption instance /// /// For example: @@ -24,7 +28,11 @@ class SortOption { /// // Sort channels by the last message date: /// final sorting = SortOption("last_message_at") /// ``` - const SortOption(this.field, {this.direction = DESC}); + const SortOption( + this.field, { + this.direction = DESC, + this.comparator, + }); /// Serialize model to json Map toJson() => _$SortOptionToJson(this); @@ -67,7 +75,7 @@ class PaginationParams { /// ``` const PaginationParams({ this.limit = 10, - this.offset, + this.offset = 0, this.greaterThan, this.greaterThanOrEqual, this.lessThan, diff --git a/packages/stream_chat/lib/src/client.dart b/packages/stream_chat/lib/src/client.dart index a1302539..55212b89 100644 --- a/packages/stream_chat/lib/src/client.dart +++ b/packages/stream_chat/lib/src/client.dart @@ -9,6 +9,7 @@ import 'package:rxdart/rxdart.dart'; import 'package:stream_chat/src/api/retry_policy.dart'; import 'package:stream_chat/src/event_type.dart'; import 'package:stream_chat/src/models/attachment_file.dart'; +import 'package:stream_chat/src/models/channel_model.dart'; import 'package:stream_chat/src/models/own_user.dart'; import 'package:stream_chat/src/platform_detector/platform_detector.dart'; import 'package:stream_chat/version.dart'; @@ -22,6 +23,7 @@ import 'api/responses.dart'; import 'api/websocket.dart'; import 'db/chat_persistence_client.dart'; import 'exceptions.dart'; +import 'models/channel_state.dart'; import 'models/event.dart'; import 'models/message.dart'; import 'models/user.dart'; @@ -489,7 +491,7 @@ class StreamChatClient { if (status == ConnectionStatus.connected && state.channels?.isNotEmpty == true) { - unawaited(queryChannels(filter: { + unawaited(_queryChannels(filter: { 'cid': { '\$in': state.channels.keys.toList(), }, @@ -578,15 +580,15 @@ class StreamChatClient { final _queryChannelsStreams = >>{}; /// Requests channels with a given query. - Future> queryChannels({ + Stream> queryChannels({ Map filter, - List sort, + List> sort, Map options, PaginationParams paginationParams = const PaginationParams(limit: 10), int messageLimit, - bool onlyOffline = false, + bool preferOffline = true, bool waitForConnect = true, - }) async { + }) async* { if (waitForConnect) { if (_connectCompleter != null && !_connectCompleter.isCompleted) { logger.info('awaiting connection completer'); @@ -598,7 +600,7 @@ class StreamChatClient { if (persistenceEnabled) { logger.warning( '$errorMessage\nTrying to retrieve channels from the offline storage.'); - onlyOffline = true; + preferOffline = true; } else { throw Exception(errorMessage); } @@ -606,34 +608,41 @@ class StreamChatClient { } final hash = base64.encode(utf8.encode( - '$filter${_asMap(sort)}$options${paginationParams?.toJson()}$messageLimit$onlyOffline')); + '$filter${_asMap(sort)}$options${paginationParams?.toJson()}$messageLimit$preferOffline', + )); + if (_queryChannelsStreams.containsKey(hash)) { - return _queryChannelsStreams[hash]; + yield await _queryChannelsStreams[hash]; + } else { + if (true) { + final channels = await _queryChannelsOffline( + filter: filter, + sort: sort, + paginationParams: paginationParams, + ); + if (channels.isNotEmpty) yield channels; + } + + final newQueryChannelsFuture = _queryChannels( + filter: filter, + sort: sort, + options: options, + paginationParams: paginationParams, + messageLimit: messageLimit, + ); + + _queryChannelsStreams[hash] = newQueryChannelsFuture; + + yield await newQueryChannelsFuture; } - - final newQueryChannelsStream = _doQueryChannels( - filter: filter, - sort: sort, - options: options, - paginationParams: paginationParams, - messageLimit: messageLimit, - onlyOffline: onlyOffline, - ).whenComplete(() { - _queryChannelsStreams.remove(hash); - }); - - _queryChannelsStreams[hash] = newQueryChannelsStream; - - return newQueryChannelsStream; } - Future> _doQueryChannels({ + Future> _queryChannels({ @required Map filter, - @required List sort, - @required Map options, - @required int messageLimit, + List> sort, + Map options, + int messageLimit, PaginationParams paginationParams = const PaginationParams(limit: 10), - bool onlyOffline = false, }) async { logger.info('Query channel start'); final defaultOptions = { @@ -661,86 +670,85 @@ class StreamChatClient { payload.addAll(paginationParams.toJson()); } - if (onlyOffline) { - return _queryChannelsOffline( - filter: filter, - sort: sort, - paginationParams: paginationParams, - ); - } + final response = await get( + '/channels', + queryParameters: { + 'payload': jsonEncode(payload), + }, + ); - try { - final response = await get( - '/channels', - queryParameters: { - 'payload': jsonEncode(payload), - }, - ); + final res = decode( + response.data, + QueryChannelsResponse.fromJson, + ); - final res = decode( - response.data, - QueryChannelsResponse.fromJson, - ); - - final users = res.channels - ?.expand((channel) => channel.members.map((member) => member.user)) - ?.toList(); - - if (users != null) { - state._updateUsers(users); - } - - logger.info('Got ${res.channels?.length} channels from api'); - - if (res.channels?.isEmpty != false && - (paginationParams?.offset ?? 0) == 0) { - logger.warning('''We could not find any channel for this query. + if (res.channels?.isEmpty == true && (paginationParams?.offset ?? 0) == 0) { + logger.warning('''We could not find any channel for this query. Please make sure to take a look at the Flutter tutorial: https://getstream.io/chat/flutter/tutorial If your application already has users and channels, you might need to adjust your query channel as explained in the docs https://getstream.io/chat/docs/query_channels/?language=dart'''); - } - - final newChannels = Map.from(state.channels ?? {}); - final channels = []; - - if (res.channels != null) { - for (final channelState in res.channels) { - final channel = newChannels[channelState.channel.cid]; - if (channel != null) { - channel.state?.updateChannelState(channelState); - channels.add(channel); - } else { - final newChannel = Channel.fromState(this, channelState); - await chatPersistenceClient - ?.updateChannelState(newChannel.state.channelState); - newChannel.state?.updateChannelState(channelState); - newChannels[newChannel.cid] = newChannel; - channels.add(newChannel); - } - } - } - - state.channels = newChannels; - - await chatPersistenceClient?.updateChannelQueries( - filter, - res.channels.map((c) => c.channel.cid).toList(), - paginationParams?.offset == null || paginationParams.offset == 0, - ); - - return channels; - } catch (e) { - if (!persistenceEnabled) { - rethrow; - } - return _queryChannelsOffline( - filter: filter, - sort: sort, - paginationParams: paginationParams, - ); + return []; } + + final channels = res.channels; + + final users = channels + .expand((it) => it.members) + .map((it) => it.user) + .toList(growable: false); + + state._updateUsers(users); + + logger.info('Got ${res.channels?.length} channels from api'); + + final updateData = _mapChannelStateToChannel(channels); + + await chatPersistenceClient?.updateChannelQueries( + filter, + channels.map((c) => c.channel.cid).toList(), + paginationParams?.offset == null || paginationParams.offset == 0, + ); + + state.channels = updateData.key; + return updateData.value; } - dynamic _parseError(DioError error) { + Future> _queryChannelsOffline({ + @required Map filter, + @required List> sort, + PaginationParams paginationParams = const PaginationParams(limit: 10), + }) async { + final offlineChannels = await chatPersistenceClient?.getChannelStates( + filter: filter, + sort: sort, + paginationParams: paginationParams, + ); + final updatedData = _mapChannelStateToChannel(offlineChannels); + state.channels = updatedData.key; + return updatedData.value; + } + + MapEntry, List> _mapChannelStateToChannel( + List channelStates, + ) { + final channels = {...state.channels ?? {}}; + final newChannels = []; + if (channelStates != null) { + for (final channelState in channelStates) { + final channel = channels[channelState.channel.cid]; + if (channel != null) { + channel.state?.updateChannelState(channelState); + newChannels.add(channel); + } else { + final newChannel = Channel.fromState(this, channelState); + channels[newChannel.cid] = newChannel; + newChannels.add(newChannel); + } + } + } + return MapEntry(channels, newChannels); + } + + Object _parseError(DioError error) { if (error.type == DioErrorType.RESPONSE) { final apiError = ApiError(error.response?.data, error.response?.statusCode); @@ -751,39 +759,6 @@ class StreamChatClient { return error; } - Future> _queryChannelsOffline({ - @required Map filter, - @required List sort, - PaginationParams paginationParams = const PaginationParams(limit: 10), - }) async { - final offlineChannels = await chatPersistenceClient?.getChannelStates( - filter: filter, - sort: sort, - paginationParams: paginationParams, - ) ?? - []; - final newChannels = Map.from(state.channels ?? {}); - logger.info('Got ${offlineChannels.length} channels from storage'); - final channels = offlineChannels.map((channelState) { - final channel = newChannels[channelState.channel.cid]; - if (channel != null) { - channel.state?.updateChannelState(channelState); - return channel; - } else { - final newChannel = Channel.fromState(this, channelState); - chatPersistenceClient - ?.updateChannelState(newChannel.state.channelState); - newChannels[newChannel.cid] = newChannel; - return newChannel; - } - }).toList(); - - if (channels.isNotEmpty) { - state.channels = newChannels; - } - return channels; - } - /// Handy method to make http GET request with error parsing. Future> get( String path, { @@ -1448,18 +1423,16 @@ class ClientState { _userController.add(user); } - void _updateUsers(List users) { - users?.forEach(_updateUser); - } - - void _updateUser(User user) { + void _updateUsers(List userList) { final newUsers = { ...users ?? {}, - user.id: user, + for (var user in userList) user.id: user, }; _usersController.add(newUsers); } + void _updateUser(User user) => _updateUsers([user]); + /// The current user OwnUser get user => _userController.value; diff --git a/packages/stream_chat/lib/src/db/chat_persistence_client.dart b/packages/stream_chat/lib/src/db/chat_persistence_client.dart index 9c31eca9..e95dadb2 100644 --- a/packages/stream_chat/lib/src/db/chat_persistence_client.dart +++ b/packages/stream_chat/lib/src/db/chat_persistence_client.dart @@ -68,23 +68,19 @@ abstract class ChatPersistenceClient { PaginationParams messagePagination, PaginationParams pinnedMessagePagination, }) async { - final members = await getMembersByCid(cid); - final reads = await getReadsByCid(cid); - final channel = await getChannelByCid(cid); - final messages = await getMessagesByCid( - cid, - messagePagination: messagePagination, - ); - final pinnedMessages = await getPinnedMessagesByCid( - cid, - messagePagination: pinnedMessagePagination, - ); + final data = await Future.wait([ + getMembersByCid(cid), + getReadsByCid(cid), + getChannelByCid(cid), + getMessagesByCid(cid, messagePagination: messagePagination), + getPinnedMessagesByCid(cid,messagePagination: pinnedMessagePagination), + ]); return ChannelState( - members: members, - read: reads, - messages: messages, - pinnedMessages: pinnedMessages, - channel: channel, + members: data[0], + read: data[1], + channel: data[2], + messages: data[3], + pinnedMessages: data[4], ); } @@ -94,7 +90,7 @@ abstract class ChatPersistenceClient { /// for filtering out states. Future> getChannelStates({ Map filter, - List sort = const [], + List> sort = const [], PaginationParams paginationParams, }); diff --git a/packages/stream_chat/lib/src/models/channel_state.dart b/packages/stream_chat/lib/src/models/channel_state.dart index 36d68a0c..c2a3d48c 100644 --- a/packages/stream_chat/lib/src/models/channel_state.dart +++ b/packages/stream_chat/lib/src/models/channel_state.dart @@ -8,7 +8,7 @@ import 'message.dart'; part 'channel_state.g.dart'; -/// The class that contains the information about a command +/// The class that contains the information about a channel @JsonSerializable() class ChannelState { /// The channel to which this state belongs diff --git a/packages/stream_chat/test/version_test.dart b/packages/stream_chat/test/version_test.dart index ada2a85c..f5b2a2c0 100644 --- a/packages/stream_chat/test/version_test.dart +++ b/packages/stream_chat/test/version_test.dart @@ -1,5 +1,6 @@ import 'dart:io'; +import 'package:rxdart/rxdart.dart'; import 'package:test/test.dart'; import 'package:stream_chat/version.dart'; @@ -10,14 +11,113 @@ void prepareTest() { } } -void main() { - prepareTest(); - test('stream chat version matches pubspec', () { - final String pubspecPath = '${Directory.current.path}/pubspec.yaml'; - final String pubspec = File(pubspecPath).readAsStringSync(); - final RegExp regex = RegExp('version:\s*(.*)'); - final RegExpMatch match = regex.firstMatch(pubspec); - expect(match, isNotNull); - expect(PACKAGE_VERSION, match.group(1).trim()); - }); +// void main() { +// prepareTest(); +// test('stream chat version matches pubspec', () { +// final String pubspecPath = '${Directory.current.path}/pubspec.yaml'; +// final String pubspec = File(pubspecPath).readAsStringSync(); +// final RegExp regex = RegExp('version:\s*(.*)'); +// final RegExpMatch match = regex.firstMatch(pubspec); +// expect(match, isNotNull); +// expect(PACKAGE_VERSION, match.group(1).trim()); +// }); +// } + +void main() async { + var items = [ + ABC('Sahil', 22, 62.0), + ABC('Devraj', 23, 76.0), + ABC('Harsh', 18, 48.0), + ABC('Harsh', 17, 88.0), + ABC('Harsh', 17, 48.0), + ABC('Devraj', 23, 74.0, { + 'Test': 'Sahil', + }), + ABC('Devraj', 12, 76.0, { + 'Test': 'Avni', + }), + ]; + + var comparators = [ + (ABC a, ABC b) { + // if (a.extraData == null) return -1; + // if (b.extraData == null) return 1; + // if (a.extraData == null && b.extraData == null) return 0; + var aa = (a.extraData ?? {})['Test'] as String; + var bb = (b.extraData ?? {})['Test'] as String; + return aa.compareTo(bb); + }, + // (ABC a, ABC b) => a.name.compareTo(b.name), + // (ABC a, ABC b) => a.age.compareTo(b.age), + // (ABC a, ABC b) => a.weight.compareTo(b.weight), + ]; + + // for (var comp in comparators.reversed) { + // items.sort(comp); + // } + + // items.sort(comparators[last]) + + Stream getLaugh2() { + if (true) { + return Stream.value('HEHO'); + } else { + return getLaugh2(); + } + } + + Stream getLaugh() async* { + yield 'HAHA'; + await Future.delayed(const Duration(seconds: 3)); + yield 'HOHO'; + await Future.delayed(const Duration(seconds: 3)); + yield 'HEHE'; + } + + // final stream = BehaviorSubject.seeded('HAHA'); + // + // stream.('HOHO'); + + await for (var value in getLaugh2()) { + print(value); + } + + // stream.add('HEHE'); + + // compare(ABC a, ABC b) { + // int result; + // for (final comparator in comparators) { + // try { + // result = comparator(a, b); + // } catch (e) { + // result = 0; + // } + // if (result != 0) return result; + // } + // return 0; + // } + // + // items.sort(compare); + // + // print(items); +} + +class ABC { + final String name; + final int age; + final double weight; + final Map extraData; + + const ABC(this.name, this.age, this.weight, [this.extraData]); + + @override + String toString() { + return ''' + \n + Name : $name, + Age : $age, + Weight : $weight, + ExtraData : $extraData, + '''; + } } diff --git a/packages/stream_chat_flutter/lib/src/channel_info.dart b/packages/stream_chat_flutter/lib/src/channel_info.dart index 4244c910..53a621a0 100644 --- a/packages/stream_chat_flutter/lib/src/channel_info.dart +++ b/packages/stream_chat_flutter/lib/src/channel_info.dart @@ -49,8 +49,11 @@ class ChannelInfo extends StatelessWidget { var alternativeWidget; if (channel.memberCount != null && channel.memberCount > 2) { + var text = '${channel.memberCount} Members'; + final watcherCount = channel.state.watcherCount ?? 0; + if (watcherCount > 0) text += ' $watcherCount Online'; alternativeWidget = Text( - '${channel.memberCount} Members, ${channel.state.watcherCount} Online', + text, style: StreamChatTheme.of(context) .channelTheme .channelHeaderTheme diff --git a/packages/stream_chat_flutter/lib/src/channel_list_view.dart b/packages/stream_chat_flutter/lib/src/channel_list_view.dart index 90b40be9..7f9818a7 100644 --- a/packages/stream_chat_flutter/lib/src/channel_list_view.dart +++ b/packages/stream_chat_flutter/lib/src/channel_list_view.dart @@ -99,7 +99,7 @@ class ChannelListView extends StatefulWidget { /// Sorting is based on field and direction, multiple sorting options can be provided. /// You can sort based on last_updated, last_message_at, updated_at, created_at or member_count. /// Direction can be ascending or descending. - final List sort; + final List> sort; /// Pagination parameters /// limit: the number of channels to return (max is 30) @@ -159,54 +159,42 @@ class ChannelListView extends StatefulWidget { _ChannelListViewState createState() => _ChannelListViewState(); } -class _ChannelListViewState extends State - with WidgetsBindingObserver { - final ScrollController _scrollController = ScrollController(); - final SlidableController _slideController = SlidableController(); - final ChannelListController _channelListController = ChannelListController(); +class _ChannelListViewState extends State { + final _slideController = SlidableController(); + + final _channelListController = ChannelListController(); @override Widget build(BuildContext context) { - var child = ChannelListCore( - channelListController: _channelListController, - listBuilder: widget.listBuilder ?? - (context, list) { - return _buildListView(list); - }, - emptyBuilder: widget.emptyBuilder ?? - (BuildContext context) { - return _buildEmptyWidget(); - }, - errorBuilder: widget.errorBuilder ?? - (BuildContext context, dynamic error) { - return _buildErrorWidget(context); - }, - loadingBuilder: widget.loadingBuilder ?? - (BuildContext context) { - return _buildLoadingWidget(); - }, + Widget child = ChannelListCore( pagination: widget.pagination, options: widget.options, sort: widget.sort, filter: widget.filter, + channelListController: _channelListController, + listBuilder: widget.listBuilder ?? _buildListView, + emptyBuilder: widget.emptyBuilder ?? _buildEmptyWidget, + errorBuilder: widget.errorBuilder ?? _buildErrorWidget, + loadingBuilder: widget.loadingBuilder ?? _buildLoadingWidget, ); - if (!widget.pullToRefresh) { - return child; - } else { - return RefreshIndicator( - onRefresh: () async { - _channelListController.loadData(); - }, + if (widget.pullToRefresh) { + child = RefreshIndicator( + onRefresh: () async => _channelListController.loadData(), child: child, ); } + + return LazyLoadScrollView( + onEndOfPage: () async { + _channelListController.paginateData(); + }, + child: child, + ); } - Widget _buildListView( - List channels, - ) { - var child; + Widget _buildListView(BuildContext context, List channels) { + Widget child; if (channels.isNotEmpty) { if (widget.crossAxisCount > 1) { @@ -216,7 +204,6 @@ class _ChannelListViewState extends State crossAxisCount: widget.crossAxisCount), itemCount: channels.length, physics: AlwaysScrollableScrollPhysics(), - controller: _scrollController, itemBuilder: (context, index) { return _gridItemBuilder(context, index, channels); }, @@ -236,18 +223,17 @@ class _ChannelListViewState extends State itemBuilder: (context, index) { return _listItemBuilder(context, index, channels); }, - controller: _scrollController, ); } } return AnimatedSwitcher( child: child, - duration: Duration(milliseconds: 500), + duration: const Duration(milliseconds: 500), ); } - Widget _buildEmptyWidget() { + Widget _buildEmptyWidget(BuildContext context) { return LayoutBuilder( builder: (context, viewportConstraints) { return SingleChildScrollView( @@ -326,7 +312,7 @@ class _ChannelListViewState extends State ); } - Widget _buildLoadingWidget() { + Widget _buildLoadingWidget(BuildContext context) { return ListView( padding: widget.padding, physics: AlwaysScrollableScrollPhysics(), @@ -341,13 +327,13 @@ class _ChannelListViewState extends State return _separatorBuilder(context, i); } } - return _buildLoadingItem(); + return _buildLoadingItem(context); }, ), ); } - Shimmer _buildLoadingItem() { + Shimmer _buildLoadingItem(BuildContext context) { if (widget.crossAxisCount > 1) { return Shimmer.fromColors( baseColor: StreamChatTheme.of(context).colorTheme.greyGainsboro, @@ -443,9 +429,7 @@ class _ChannelListViewState extends State } } - Widget _buildErrorWidget( - BuildContext context, - ) { + Widget _buildErrorWidget(BuildContext context, Object error) { return Center( child: Column( mainAxisAlignment: MainAxisAlignment.center, @@ -572,23 +556,13 @@ class _ChannelListViewState extends State ], child: Container( color: StreamChatTheme.of(context).colorTheme.whiteSnow, - child: widget.channelPreviewBuilder != null - ? widget.channelPreviewBuilder( - context, - channel, - ) - : ChannelPreview( - onLongPress: widget.onChannelLongPress, - channel: channel, - onImageTap: widget.onImageTap != null - ? () { - widget.onImageTap(channel); - } - : null, - onTap: (channel) { - onTap(channel, widget.channelWidget); - }, - ), + child: widget.channelPreviewBuilder?.call(context, channel) ?? + ChannelPreview( + onLongPress: widget.onChannelLongPress, + channel: channel, + onImageTap: widget.onImageTap?.call(channel), + onTap: (channel) => onTap(channel, widget.channelWidget), + ), ), ); }, @@ -618,9 +592,7 @@ class _ChannelListViewState extends State width: 64, height: 64, ), - onTap: () { - widget.onChannelTap(channel, null); - }, + onTap: () => widget.onChannelTap(channel, null), ), SizedBox(height: 7), Padding( @@ -680,70 +652,4 @@ class _ChannelListViewState extends State color: effect.color.withOpacity(effect.alpha ?? 1.0), ); } - - void _listenChannelPagination(ChannelsBlocState channelsProvider) { - if (_scrollController.position.maxScrollExtent == - _scrollController.offset && - _scrollController.offset != 0) { - _channelListController.paginateData(); - } - } - - StreamSubscription _subscription; - - @override - void initState() { - super.initState(); - - WidgetsBinding.instance.addObserver(this); - - final channelsBloc = ChannelsBloc.of(context); - channelsBloc.queryChannels( - filter: widget.filter, - sortOptions: widget.sort, - paginationParams: widget.pagination, - options: widget.options, - ); - - _scrollController.addListener(() { - channelsBloc.queryChannelsLoading.first.then((loading) { - if (!loading) { - _listenChannelPagination(channelsBloc); - } - }); - }); - - final client = StreamChat.of(context).client; - - _subscription = client - .on( - EventType.connectionRecovered, - EventType.notificationAddedToChannel, - EventType.notificationMessageNew, - EventType.channelVisible, - ) - .listen((event) { - _channelListController.loadData(); - }); - } - - @override - void didUpdateWidget(ChannelListView oldWidget) { - super.didUpdateWidget(oldWidget); - - if (widget.filter?.toString() != oldWidget.filter?.toString() || - jsonEncode(widget.sort) != jsonEncode(oldWidget.sort) || - widget.pagination?.toJson()?.toString() != - oldWidget.pagination?.toJson()?.toString() || - widget.options?.toString() != oldWidget.options?.toString()) { - _channelListController.loadData(); - } - } - - @override - void dispose() { - _subscription.cancel(); - WidgetsBinding.instance.removeObserver(this); - super.dispose(); - } } diff --git a/packages/stream_chat_flutter_core/lib/src/channel_list_core.dart b/packages/stream_chat_flutter_core/lib/src/channel_list_core.dart index 18a8e477..acc9f802 100644 --- a/packages/stream_chat_flutter_core/lib/src/channel_list_core.dart +++ b/packages/stream_chat_flutter_core/lib/src/channel_list_core.dart @@ -106,7 +106,7 @@ class ChannelListCore extends StatefulWidget { /// Sorting is based on field and direction, multiple sorting options can be provided. /// You can sort based on last_updated, last_message_at, updated_at, created_at or member_count. /// Direction can be ascending or descending. - final List sort; + final List> sort; /// Pagination parameters /// limit: the number of channels to return (max is 30) @@ -118,8 +118,7 @@ class ChannelListCore extends StatefulWidget { _ChannelListCoreState createState() => _ChannelListCoreState(); } -class _ChannelListCoreState extends State - with WidgetsBindingObserver { +class _ChannelListCoreState extends State { @override Widget build(BuildContext context) { final channelsBloc = ChannelsBloc.of(context); @@ -133,34 +132,21 @@ class _ChannelListCoreState extends State return StreamBuilder>( stream: channelsBlocState.channelsStream, builder: (context, snapshot) { - var child; if (snapshot.hasError) { - child = _buildErrorWidget( - snapshot, - context, - channelsBlocState, - ); - } else if (!snapshot.hasData) { - child = _buildLoadingWidget(); - } else { - final channels = snapshot.data; - - child = widget.emptyBuilder(context); - - if (channels.isNotEmpty) { - return widget.listBuilder(context, channels); - } + return _buildErrorWidget(snapshot, context, channelsBlocState); } - - return child; + if (!snapshot.hasData) { + return widget.loadingBuilder(context); + } + final channels = snapshot.data; + if (channels.isEmpty) { + return widget.emptyBuilder(context); + } + return widget.listBuilder(context, channels); }, ); } - Widget _buildLoadingWidget() { - return widget.loadingBuilder(context); - } - Widget _buildErrorWidget( AsyncSnapshot> snapshot, BuildContext context, @@ -197,39 +183,21 @@ class _ChannelListCoreState extends State ); } - StreamSubscription _subscription; + StreamSubscription _subscription; @override void initState() { super.initState(); - - WidgetsBinding.instance.addObserver(this); - - final channelsBloc = ChannelsBloc.of(context); - channelsBloc.queryChannels( - filter: widget.filter, - sortOptions: widget.sort, - paginationParams: widget.pagination, - options: widget.options, - ); - + loadData(); final client = StreamChatCore.of(context).client; - _subscription = client .on( - EventType.connectionRecovered, - EventType.notificationAddedToChannel, - EventType.notificationMessageNew, - EventType.channelVisible, - ) - .listen((event) { - channelsBloc.queryChannels( - filter: widget.filter, - sortOptions: widget.sort, - paginationParams: widget.pagination, - options: widget.options, - ); - }); + EventType.connectionRecovered, + EventType.notificationAddedToChannel, + EventType.notificationMessageNew, + EventType.channelVisible, + ) + .listen((event) => loadData()); if (widget.channelListController != null) { widget.channelListController.loadData = loadData; @@ -246,20 +214,13 @@ class _ChannelListCoreState extends State widget.pagination?.toJson()?.toString() != oldWidget.pagination?.toJson()?.toString() || widget.options?.toString() != oldWidget.options?.toString()) { - final channelsBloc = ChannelsBloc.of(context); - channelsBloc.queryChannels( - filter: widget.filter, - sortOptions: widget.sort, - paginationParams: widget.pagination, - options: widget.options, - ); + loadData(); } } @override void dispose() { _subscription.cancel(); - WidgetsBinding.instance.removeObserver(this); super.dispose(); } } diff --git a/packages/stream_chat_flutter_core/lib/src/channels_bloc.dart b/packages/stream_chat_flutter_core/lib/src/channels_bloc.dart index 7009234b..836e57c8 100644 --- a/packages/stream_chat_flutter_core/lib/src/channels_bloc.dart +++ b/packages/stream_chat_flutter_core/lib/src/channels_bloc.dart @@ -84,7 +84,7 @@ class ChannelsBlocState extends State /// Calls [client.queryChannels] updating [queryChannelsLoading] stream Future queryChannels({ Map filter, - List sortOptions, + List> sortOptions, PaginationParams paginationParams, Map options, bool onlyOffline = false, @@ -102,21 +102,21 @@ class ChannelsBlocState extends State paginationParams.offset == null || paginationParams.offset == 0; final oldChannels = List.from(channels ?? []); - final _channels = await client.queryChannels( + await for (final channels in client.queryChannels( filter: filter, sort: sortOptions, options: options, paginationParams: paginationParams, - onlyOffline: onlyOffline, - ); - - if (clear) { - _channelsController.add(_channels); - } else { - final l = oldChannels + _channels; - _channelsController.add(l); + preferOffline: onlyOffline, + )) { + if (clear) { + _channelsController.add(channels); + } else { + final l = oldChannels + channels; + _channelsController.add(l); + } + _queryChannelsLoadingController.sink.add(false); } - _queryChannelsLoadingController.sink.add(false); } catch (err, stackTrace) { print(err); print(stackTrace); diff --git a/packages/stream_chat_persistence/example/pubspec.yaml b/packages/stream_chat_persistence/example/pubspec.yaml index 5008566b..2586bea6 100644 --- a/packages/stream_chat_persistence/example/pubspec.yaml +++ b/packages/stream_chat_persistence/example/pubspec.yaml @@ -11,7 +11,7 @@ dependencies: flutter: sdk: flutter cupertino_icons: ^1.0.0 - stream_chat: + stream_chat: path: ../../stream_chat stream_chat_persistence: path: ../ diff --git a/packages/stream_chat_persistence/lib/src/dao/channel_query_dao.dart b/packages/stream_chat_persistence/lib/src/dao/channel_query_dao.dart index 83743deb..7d4af560 100644 --- a/packages/stream_chat_persistence/lib/src/dao/channel_query_dao.dart +++ b/packages/stream_chat_persistence/lib/src/dao/channel_query_dao.dart @@ -15,9 +15,7 @@ part 'channel_query_dao.g.dart'; class ChannelQueryDao extends DatabaseAccessor with _$ChannelQueryDaoMixin { /// Creates a new channel query dao instance - ChannelQueryDao(this._db) : super(_db); - - final MoorChatDatabase _db; + ChannelQueryDao(MoorChatDatabase db) : super(db); String _computeHash(Map filter) { if (filter == null) { @@ -37,19 +35,19 @@ class ChannelQueryDao extends DatabaseAccessor ) async { final hash = _computeHash(filter); if (clearQueryCache) { - await (delete(channelQueries) - ..where((query) => query.queryHash.equals(hash))) - .go(); + await batch((it) { + it.deleteWhere( + channelQueries, + (c) => c.queryHash.equals(hash), + ); + }); } - return batch((batch) { - batch.insertAll( + return batch((it) { + it.insertAll( channelQueries, cids.map((cid) { - return ChannelQueryEntity( - queryHash: hash, - channelCid: cid, - ); + return ChannelQueryEntity(queryHash: hash, channelCid: cid); }).toList(), mode: InsertMode.insertOrReplace, ); @@ -57,70 +55,74 @@ class ChannelQueryDao extends DatabaseAccessor } /// Get list of channels by filter, sort and paginationParams - Future> getChannelStates({ + Future> getChannels({ Map filter, - List sort = const [], + List> sort = const [], PaginationParams paginationParams, }) async { + assert(() { + if (sort != null && sort.any((it) => it.comparator == null)) { + throw ArgumentError( + 'SortOption requires a comparator in order to sort', + ); + } + return true; + }()); + final hash = _computeHash(filter); - final cachedChannels = await Future.wait(await (select(channelQueries) + final cachedChannelCids = await (select(channelQueries) ..where((c) => c.queryHash.equals(hash))) - .get() - .then((channelQueries) { - final cids = channelQueries.map((c) => c.channelCid).toList(); - final query = select(channels)..where((c) => c.cid.isIn(cids)); + .map((c) => c.channelCid) + .get(); - sort = sort - ?.where((s) => ChannelModel.topLevelFields.contains(s.field)) - ?.toList(); + final query = select(channels)..where((c) => c.cid.isIn(cachedChannelCids)); - if (sort != null && sort.isNotEmpty) { - query.orderBy(sort.map((s) { - final orderExpression = CustomExpression('channels.${s.field}'); - return (c) => OrderingTerm( - expression: orderExpression, - mode: s.direction == 1 ? OrderingMode.asc : OrderingMode.desc, - ); - }).toList()); - } + final cachedChannels = await (query.join([ + leftOuterJoin(users, channels.createdById.equalsExp(users.id)), + ]).map((row) { + final createdByEntity = row.readTable(users); + final channelEntity = row.readTable(channels); + return channelEntity.toChannelModel(createdBy: createdByEntity?.toUser()); + })).get(); - if (paginationParams != null) { - query.limit( - paginationParams.limit ?? 10, - offset: paginationParams.offset, - ); - } + final possibleSortingFields = cachedChannels.fold>( + ChannelModel.topLevelFields, (previousValue, element) { + return {...previousValue, ...element.extraData.keys}.toList(); + }); - return query.join([ - leftOuterJoin(users, channels.createdById.equalsExp(users.id)), - ]).map((row) async { - final userEntity = row.readTable(users); - final channelEntity = row.readTable(channels); + sort = sort + ?.where((s) => possibleSortingFields.contains(s.field)) + ?.toList(growable: false); - final cid = channelEntity.cid; - final members = await _db.memberDao.getMembersByCid(cid); - final reads = await _db.readDao.getReadsByCid(cid); - final messages = await _db.messageDao.getMessagesByCid(cid); - final pinnedMessages = await _db.pinnedMessageDao.getMessagesByCid(cid); + Comparator chainedComparator = (a, b) { + final dateA = a.lastMessageAt ?? a.createdAt; + final dateB = b.lastMessageAt ?? b.createdAt; + return dateB.compareTo(dateA); + }; - return channelEntity.toChannelState( - createdBy: userEntity?.toUser(), - members: members, - reads: reads, - messages: messages, - pinnedMessages: pinnedMessages, - ); - }).get(); - })); + if (sort != null && sort.isNotEmpty) { + chainedComparator = (a, b) { + int result; + for (final comparator in sort.map((it) => it.comparator)) { + try { + result = comparator(a, b); + } catch (e) { + result = 0; + } + if (result != 0) return result; + } + return 0; + }; + } - if (sort?.isEmpty != false && cachedChannels?.isNotEmpty == true) { - cachedChannels - .sort((a, b) => b.channel.updatedAt.compareTo(a.channel.updatedAt)); - cachedChannels.sort((a, b) { - final dateA = a.channel.lastMessageAt ?? a.channel.createdAt; - final dateB = b.channel.lastMessageAt ?? b.channel.createdAt; - return dateB.compareTo(dateA); - }); + cachedChannels.sort(chainedComparator); + + if (paginationParams?.offset != null) { + cachedChannels.removeRange(0, paginationParams.offset); + } + + if (paginationParams?.limit != null) { + return cachedChannels.take(paginationParams.limit).toList(); } return cachedChannels; diff --git a/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart b/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart index 50d17111..e6be01ec 100644 --- a/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart +++ b/packages/stream_chat_persistence/lib/src/stream_chat_persistence_client.dart @@ -161,14 +161,15 @@ class StreamChatPersistenceClient extends ChatPersistenceClient { @override Future> getChannelStates({ Map filter, - List sort = const [], + List> sort = const [], PaginationParams paginationParams, - }) { - return _db.channelQueryDao.getChannelStates( + }) async { + final channels = await _db.channelQueryDao.getChannels( filter: filter, sort: sort, paginationParams: paginationParams, ); + return Future.wait(channels.map((e) => getChannelStateByCid(e.cid))); } @override