import 'dart:async'; import 'package:flutter/material.dart'; import 'package:rxdart/rxdart.dart'; import 'package:stream_chat/stream_chat.dart'; enum QueryDirection { top, bottom } /// Widget used to provide information about the channel to the widget tree /// /// Use [StreamChannel.of] to get the current [StreamChannelState] instance. class StreamChannel extends StatefulWidget { const StreamChannel({ Key key, @required this.child, @required this.channel, this.showLoading = true, this.initialMessageId, }) : assert(child != null), assert(channel != null), super(key: key); final Widget child; final Channel channel; final bool showLoading; /// If passed the channel will load from this particular message. final String initialMessageId; /// Use this method to get the current [StreamChannelState] instance static StreamChannelState of(BuildContext context) { StreamChannelState streamChannelState; streamChannelState = context.findAncestorStateOfType(); if (streamChannelState == null) { throw Exception( 'You must have a StreamChannel widget at the top of your widget tree'); } return streamChannelState; } @override StreamChannelState createState() => StreamChannelState(); } class StreamChannelState extends State { /// Current channel Channel get channel => widget.channel; /// InitialMessageId String get initialMessageId => widget.initialMessageId; /// Current channel state stream Stream get channelStateStream => widget.channel.state.channelStateStream; final _queryTopMessagesController = BehaviorSubject.seeded(false); final _queryBottomMessagesController = BehaviorSubject.seeded(false); /// The stream notifying the state of [_queryTopMessages] call Stream get queryTopMessages => _queryTopMessagesController.stream; /// The stream notifying the state of [_queryBottomMessages] call Stream get queryBottomMessages => _queryBottomMessagesController.stream; bool _topPaginationEnded = false; bool _bottomPaginationEnded = false; Future _queryTopMessages({ int limit = 20, bool preferOffline = false, }) async { if (_topPaginationEnded || _queryTopMessagesController?.value == true) { return; } _queryTopMessagesController.add(true); if (channel.state.messages.isEmpty) { return _queryTopMessagesController.add(false); } final oldestMessage = channel.state.messages.first; try { final state = await queryBeforeMessage( oldestMessage.id, limit: limit, preferOffline: preferOffline, ); if (state.messages.isEmpty || state.messages.length < limit) { _topPaginationEnded = true; } _queryTopMessagesController.add(false); } catch (e, stk) { _queryTopMessagesController.addError(e, stk); } } Future _queryBottomMessages({ int limit = 20, bool preferOffline = false, }) async { if (_bottomPaginationEnded || _queryBottomMessagesController?.value == true || channel?.state?.isUpToDate == true) return; _queryBottomMessagesController.add(true); if (channel.state.messages.isEmpty) { return _queryBottomMessagesController.add(false); } final recentMessage = channel.state.messages.last; try { final state = await queryAfterMessage( recentMessage.id, limit: limit, preferOffline: preferOffline, ); if (state.messages.isEmpty || state.messages.length < limit) { _bottomPaginationEnded = true; } _queryBottomMessagesController.add(false); } catch (e, stk) { _queryBottomMessagesController.addError(e, stk); } } /// Calls [channel.query] updating [queryMessage] stream Future queryMessages({QueryDirection direction = QueryDirection.top}) { if (direction == QueryDirection.top) return _queryTopMessages(); return _queryBottomMessages(); } /// Calls [channel.getReplies] updating [queryMessage] stream Future getReplies( String parentId, { int limit = 50, bool preferOffline = false, }) async { if (_topPaginationEnded || _queryTopMessagesController.value) return; _queryTopMessagesController.add(true); if (!channel.state.threads.containsKey(parentId)) { return _queryTopMessagesController.add(false); } final thread = channel.state.threads[parentId]; if (thread.isEmpty) return _queryTopMessagesController.add(false); final message = thread.first; try { final state = await queryBeforeMessage( message.id, limit: limit, preferOffline: preferOffline, ); if (state.messages.isEmpty || state.messages.length < limit) { _topPaginationEnded = true; } _queryTopMessagesController.add(false); } catch (e, stk) { _queryTopMessagesController.addError(e, stk); } } /// Query the channel members and watchers Future queryMembersAndWatchers() async { await widget.channel.query( membersPagination: PaginationParams( offset: channel.state.members?.length, limit: 100, ), watchersPagination: PaginationParams( offset: channel.state.watchers?.length, limit: 100, ), ); } /// Loads channel at specific message Future loadChannelAtMessage( String messageId, { int before = 20, int after = 20, bool preferOffline = false, }) { return queryAtMessage( messageId: messageId, before: before, after: after, preferOffline: preferOffline, ); } /// Future queryAtMessage({ String messageId, int before = 20, int after = 20, bool preferOffline = false, }) async { if (channel.state == null) return; channel.state.isUpToDate = false; channel.state.truncate(); if (messageId == null) { await channel.query( messagesPagination: PaginationParams( limit: before, ), preferOffline: preferOffline, ); channel.state.isUpToDate = true; return; } return Future.wait([ queryBeforeMessage( messageId, limit: before, preferOffline: preferOffline, ), queryAfterMessage( messageId, limit: after, preferOffline: preferOffline, ), ]); } /// Future queryBeforeMessage( String messageId, { int limit = 20, bool preferOffline = false, }) { return channel.query( messagesPagination: PaginationParams( lessThan: messageId, limit: limit, ), preferOffline: preferOffline, ); } /// Future queryAfterMessage( String messageId, { int limit = 20, bool preferOffline = false, }) async { final state = await channel.query( messagesPagination: PaginationParams( greaterThanOrEqual: messageId, limit: limit, ), preferOffline: preferOffline, ); if (state.messages.isEmpty || state.messages.length < limit) { channel.state.isUpToDate = true; } return state; } /// Reloads the channel with latest message Future reloadChannel() => queryAtMessage(before: 30); List> _futures; Future get _loadChannelAtMessage async { try { await loadChannelAtMessage(initialMessageId); return true; } catch (e, stk) { print('Error: $e\nStack: $stk'); rethrow; } } @override void initState() { super.initState(); _futures = [widget.channel.initialized]; if (initialMessageId != null) { _futures.add(_loadChannelAtMessage); } } @override void dispose() { _queryTopMessagesController.close(); _queryBottomMessagesController.close(); super.dispose(); } @override Widget build(BuildContext context) { Widget child = FutureBuilder>( future: Future.wait(_futures), initialData: [ channel.state != null, if (initialMessageId != null) false, ], builder: (context, snapshot) { if (snapshot.hasError) { if (snapshot.error is Error) { print((snapshot.error as Error).stackTrace); } var message = snapshot.error.toString(); if (snapshot.error is DioError) { final dioError = snapshot.error as DioError; if (dioError.type == DioErrorType.RESPONSE) { message = dioError.message; } else { message = 'Check your connection and retry'; } } return Center( child: Text(message), ); } final initialized = snapshot.data[0]; final dataLoaded = initialMessageId == null ? true : snapshot.data[1]; if (widget.showLoading && (!initialized || !dataLoaded)) { return Center( child: CircularProgressIndicator(), ); } return widget.child; }, ); if (initialMessageId != null) { child = Material(child: child); } return child; } }