Merge branch 'develop' into fix/channel-leave-cache
This commit is contained in:
@@ -1379,13 +1379,14 @@ class Channel {
|
||||
this.state?.updateChannelState(updatedState);
|
||||
return updatedState;
|
||||
} catch (e) {
|
||||
if (!_client.persistenceEnabled) {
|
||||
rethrow;
|
||||
if (_client.persistenceEnabled) {
|
||||
return _client.chatPersistenceClient!.getChannelStateByCid(
|
||||
cid!,
|
||||
messagePagination: messagesPagination,
|
||||
);
|
||||
}
|
||||
return _client.chatPersistenceClient!.getChannelStateByCid(
|
||||
cid!,
|
||||
messagePagination: messagesPagination,
|
||||
);
|
||||
|
||||
rethrow;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1841,9 +1842,7 @@ class ChannelClientState {
|
||||
|
||||
/// [isUpToDate] flag count as a stream.
|
||||
Stream<bool> get isUpToDateStream => _isUpToDateController.stream;
|
||||
|
||||
final BehaviorSubject<bool> _isUpToDateController =
|
||||
BehaviorSubject.seeded(true);
|
||||
final _isUpToDateController = BehaviorSubject.seeded(true);
|
||||
|
||||
/// The retry queue associated to this channel.
|
||||
late final RetryQueue _retryQueue;
|
||||
|
||||
@@ -125,10 +125,6 @@ class StreamChatClient {
|
||||
final _tokenManager = TokenManager();
|
||||
final _connectionIdManager = ConnectionIdManager();
|
||||
|
||||
set chatPersistenceClient(ChatPersistenceClient? value) {
|
||||
_originalChatPersistenceClient = value;
|
||||
}
|
||||
|
||||
/// Default user agent for all requests
|
||||
static String defaultUserAgent =
|
||||
'stream-chat-dart-client-${CurrentPlatform.name}';
|
||||
@@ -139,15 +135,15 @@ class StreamChatClient {
|
||||
/// The current package version
|
||||
static const packageVersion = PACKAGE_VERSION;
|
||||
|
||||
ChatPersistenceClient? _originalChatPersistenceClient;
|
||||
|
||||
/// Chat persistence client
|
||||
ChatPersistenceClient? get chatPersistenceClient => _chatPersistenceClient;
|
||||
ChatPersistenceClient? chatPersistenceClient;
|
||||
|
||||
ChatPersistenceClient? _chatPersistenceClient;
|
||||
|
||||
/// Whether the chat persistence is available or not
|
||||
bool get persistenceEnabled => _chatPersistenceClient != null;
|
||||
/// Returns `True` if the [chatPersistenceClient] is available and connected.
|
||||
/// Otherwise, returns `False`.
|
||||
bool get persistenceEnabled {
|
||||
final client = chatPersistenceClient;
|
||||
return client != null && client.isConnected;
|
||||
}
|
||||
|
||||
late final RetryPolicy _retryPolicy;
|
||||
|
||||
@@ -324,20 +320,27 @@ class StreamChatClient {
|
||||
final ownUser = OwnUser.fromUser(user);
|
||||
state.currentUser = ownUser;
|
||||
|
||||
if (!connectWebSocket) return ownUser;
|
||||
|
||||
try {
|
||||
if (_originalChatPersistenceClient != null) {
|
||||
_chatPersistenceClient = _originalChatPersistenceClient;
|
||||
await _chatPersistenceClient!.connect(ownUser.id);
|
||||
// Connect to persistence client if its set.
|
||||
if (chatPersistenceClient != null) {
|
||||
await openPersistenceConnection(ownUser);
|
||||
}
|
||||
final connectedUser = await openConnection(
|
||||
includeUserDetailsInConnectCall: true,
|
||||
);
|
||||
return state.currentUser = connectedUser;
|
||||
|
||||
// Connect to websocket if [connectWebSocket] is true.
|
||||
//
|
||||
// This is useful when you want to connect to websocket
|
||||
// at a later stage or use the client in connection-less mode.
|
||||
if (connectWebSocket) {
|
||||
final connectedUser = await openConnection(
|
||||
includeUserDetailsInConnectCall: true,
|
||||
);
|
||||
state.currentUser = connectedUser;
|
||||
}
|
||||
|
||||
return state.currentUser!;
|
||||
} catch (e, stk) {
|
||||
if (e is StreamWebSocketError && e.isRetriable) {
|
||||
final event = await _chatPersistenceClient?.getConnectionInfo();
|
||||
final event = await chatPersistenceClient?.getConnectionInfo();
|
||||
if (event != null) return ownUser.merge(event.me);
|
||||
}
|
||||
logger.severe('error connecting user : ${ownUser.id}', e, stk);
|
||||
@@ -345,6 +348,40 @@ class StreamChatClient {
|
||||
}
|
||||
}
|
||||
|
||||
/// Connects the [chatPersistenceClient] to the given [user].
|
||||
Future<void> openPersistenceConnection(User user) async {
|
||||
final client = chatPersistenceClient;
|
||||
if (client == null) {
|
||||
throw const StreamChatError('Chat persistence client is not set');
|
||||
}
|
||||
|
||||
if (client.isConnected) {
|
||||
// If the persistence client is already connected to the userId,
|
||||
// we don't need to connect again.
|
||||
if (client.userId == user.id) return;
|
||||
|
||||
throw const StreamChatError('''
|
||||
Chat persistence client is already connected to a different user,
|
||||
please close the connection before connecting a new one.''');
|
||||
}
|
||||
|
||||
// Connect the persistence client to the userId.
|
||||
return client.connect(user.id);
|
||||
}
|
||||
|
||||
/// Disconnects the [chatPersistenceClient] from the current user.
|
||||
Future<void> closePersistenceConnection({bool flush = false}) async {
|
||||
final client = chatPersistenceClient;
|
||||
// If the persistence client is never connected, we don't need to close it.
|
||||
if (client == null || !client.isConnected) {
|
||||
logger.info('Chat persistence client is not connected');
|
||||
return;
|
||||
}
|
||||
|
||||
// Disconnect the persistence client.
|
||||
return client.disconnect(flush: flush);
|
||||
}
|
||||
|
||||
/// Creates a new WebSocket connection with the current user.
|
||||
/// If [includeUserDetailsInConnectCall] is true it will include the current
|
||||
/// user details in the connect call.
|
||||
@@ -422,7 +459,7 @@ class StreamChatClient {
|
||||
final connectionId = event.connectionId;
|
||||
if (connectionId != null) {
|
||||
_connectionIdManager.setConnectionId(connectionId);
|
||||
_chatPersistenceClient?.updateConnectionInfo(event);
|
||||
chatPersistenceClient?.updateConnectionInfo(event);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -460,9 +497,9 @@ class StreamChatClient {
|
||||
// channels are empty, assuming it's a fresh start
|
||||
// and making sure `lastSyncAt` is initialized
|
||||
if (persistenceEnabled) {
|
||||
final lastSyncAt = await _chatPersistenceClient?.getLastSyncAt();
|
||||
final lastSyncAt = await chatPersistenceClient?.getLastSyncAt();
|
||||
if (lastSyncAt == null) {
|
||||
await _chatPersistenceClient?.updateLastSyncAt(DateTime.now());
|
||||
await chatPersistenceClient?.updateLastSyncAt(DateTime.now());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -493,13 +530,12 @@ class StreamChatClient {
|
||||
/// Will automatically fetch [cids] and [lastSyncedAt] if [persistenceEnabled]
|
||||
Future<void> sync({List<String>? cids, DateTime? lastSyncAt}) {
|
||||
return synchronized(() async {
|
||||
final channels = cids ?? await _chatPersistenceClient?.getChannelCids();
|
||||
final channels = cids ?? await chatPersistenceClient?.getChannelCids();
|
||||
if (channels == null || channels.isEmpty) {
|
||||
return;
|
||||
}
|
||||
|
||||
final syncAt =
|
||||
lastSyncAt ?? await _chatPersistenceClient?.getLastSyncAt();
|
||||
final syncAt = lastSyncAt ?? await chatPersistenceClient?.getLastSyncAt();
|
||||
if (syncAt == null) {
|
||||
return;
|
||||
}
|
||||
@@ -520,7 +556,7 @@ class StreamChatClient {
|
||||
|
||||
final now = DateTime.now();
|
||||
_lastSyncedAt = now;
|
||||
_chatPersistenceClient?.updateLastSyncAt(now);
|
||||
chatPersistenceClient?.updateLastSyncAt(now);
|
||||
} catch (e, stk) {
|
||||
logger.severe('Error during sync', e, stk);
|
||||
}
|
||||
@@ -679,7 +715,7 @@ class StreamChatClient {
|
||||
|
||||
final updateData = _mapChannelStateToChannel(channels);
|
||||
|
||||
await _chatPersistenceClient?.updateChannelQueries(
|
||||
await chatPersistenceClient?.updateChannelQueries(
|
||||
filter,
|
||||
channels.map((c) => c.channel!.cid).toList(),
|
||||
clearQueryCache: paginationParams.offset == 0,
|
||||
@@ -698,7 +734,7 @@ class StreamChatClient {
|
||||
List<SortOption<ChannelState>>? channelStateSort,
|
||||
PaginationParams paginationParams = const PaginationParams(),
|
||||
}) async {
|
||||
final offlineChannels = (await _chatPersistenceClient?.getChannelStates(
|
||||
final offlineChannels = (await chatPersistenceClient?.getChannelStates(
|
||||
filter: filter,
|
||||
// ignore: deprecated_member_use_from_same_package
|
||||
sort: sort,
|
||||
@@ -1362,7 +1398,7 @@ class StreamChatClient {
|
||||
final response =
|
||||
await _chatApi.message.deleteMessage(messageId, hard: hard);
|
||||
if (hard == true) {
|
||||
await _chatPersistenceClient?.deleteMessageById(messageId);
|
||||
await chatPersistenceClient?.deleteMessageById(messageId);
|
||||
}
|
||||
return response;
|
||||
}
|
||||
@@ -1468,34 +1504,33 @@ class StreamChatClient {
|
||||
Future<void> disconnectUser({bool flushChatPersistence = false}) async {
|
||||
logger.info('Disconnecting user : ${state.currentUser?.id}');
|
||||
|
||||
// resetting state
|
||||
// resetting state.
|
||||
state.dispose();
|
||||
state = ClientState(this);
|
||||
_lastSyncedAt = null;
|
||||
|
||||
// resetting credentials
|
||||
// resetting credentials.
|
||||
_tokenManager.reset();
|
||||
_connectionIdManager.reset();
|
||||
|
||||
// disconnecting persistence client
|
||||
await _chatPersistenceClient?.disconnect(flush: flushChatPersistence);
|
||||
_chatPersistenceClient = null;
|
||||
// closing persistence connection.
|
||||
await closePersistenceConnection(flush: flushChatPersistence);
|
||||
|
||||
// closing web-socket connection
|
||||
closeConnection();
|
||||
return closeConnection();
|
||||
}
|
||||
|
||||
/// Call this function to dispose the client
|
||||
Future<void> dispose() async {
|
||||
logger.info('Disposing new StreamChatClient');
|
||||
|
||||
// disposing state
|
||||
// disposing state.
|
||||
state.dispose();
|
||||
|
||||
// disconnecting persistence client
|
||||
await _chatPersistenceClient?.disconnect();
|
||||
// closing persistence connection.
|
||||
await closePersistenceConnection();
|
||||
|
||||
// closing web-socket connection
|
||||
// closing web-socket connection.
|
||||
closeConnection();
|
||||
|
||||
await _eventController.close();
|
||||
|
||||
@@ -89,8 +89,13 @@ class StreamChatNetworkError extends StreamChatError {
|
||||
}) : super(message);
|
||||
|
||||
///
|
||||
factory StreamChatNetworkError.fromDioError(DioError error) {
|
||||
final response = error.response;
|
||||
@Deprecated('Use `StreamChatNetworkError.fromDioException` instead')
|
||||
factory StreamChatNetworkError.fromDioError(DioException error) =
|
||||
StreamChatNetworkError.fromDioException;
|
||||
|
||||
///
|
||||
factory StreamChatNetworkError.fromDioException(DioException exception) {
|
||||
final response = exception.response;
|
||||
ErrorResponse? errorResponse;
|
||||
final data = response?.data;
|
||||
if (data != null) {
|
||||
@@ -100,12 +105,12 @@ class StreamChatNetworkError extends StreamChatError {
|
||||
code: errorResponse?.code ?? -1,
|
||||
message: errorResponse?.message ??
|
||||
response?.statusMessage ??
|
||||
error.message ??
|
||||
exception.message ??
|
||||
'',
|
||||
statusCode: errorResponse?.statusCode ?? response?.statusCode,
|
||||
data: errorResponse,
|
||||
isRequestCancelledError: error.type == DioErrorType.cancel,
|
||||
)..stackTrace = error.stackTrace;
|
||||
isRequestCancelledError: exception.type == DioExceptionType.cancel,
|
||||
)..stackTrace = exception.stackTrace;
|
||||
}
|
||||
|
||||
/// Error code
|
||||
|
||||
@@ -46,26 +46,26 @@ class AuthInterceptor extends QueuedInterceptor {
|
||||
|
||||
@override
|
||||
void onError(
|
||||
DioError err,
|
||||
DioException exception,
|
||||
ErrorInterceptorHandler handler,
|
||||
) async {
|
||||
final data = err.response?.data;
|
||||
final data = exception.response?.data;
|
||||
if (data == null || data is! Map<String, dynamic>) {
|
||||
return handler.next(err);
|
||||
return handler.next(exception);
|
||||
}
|
||||
|
||||
final error = ErrorResponse.fromJson(data);
|
||||
if (error.code == ChatErrorCode.tokenExpired.code) {
|
||||
if (_tokenManager.isStatic) return handler.next(err);
|
||||
if (_tokenManager.isStatic) return handler.next(exception);
|
||||
await _tokenManager.loadToken(refresh: true);
|
||||
try {
|
||||
final options = err.requestOptions;
|
||||
final options = exception.requestOptions;
|
||||
final response = await _client.fetch(options);
|
||||
return handler.resolve(response);
|
||||
} on DioError catch (error) {
|
||||
return handler.next(error);
|
||||
} on DioException catch (exception) {
|
||||
return handler.next(exception);
|
||||
}
|
||||
}
|
||||
return handler.next(err);
|
||||
return handler.next(exception);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -119,32 +119,32 @@ class LoggingInterceptor extends Interceptor {
|
||||
}
|
||||
|
||||
@override
|
||||
void onError(DioError err, ErrorInterceptorHandler handler) {
|
||||
void onError(DioException exception, ErrorInterceptorHandler handler) {
|
||||
if (error) {
|
||||
if (err.type == DioErrorType.badResponse) {
|
||||
final uri = err.response?.requestOptions.uri;
|
||||
if (exception.type == DioExceptionType.badResponse) {
|
||||
final uri = exception.response?.requestOptions.uri;
|
||||
_printBoxed(
|
||||
_logPrintError,
|
||||
header:
|
||||
'DioError ║ Status: ${err.response?.statusCode} ${err.response?.statusMessage}',
|
||||
'DioException ║ Status: ${exception.response?.statusCode} ${exception.response?.statusMessage}',
|
||||
text: uri.toString(),
|
||||
);
|
||||
if (err.response != null && err.response?.data != null) {
|
||||
_logPrintError('╔ ${err.type.toString()}');
|
||||
_printResponse(_logPrintError, err.response!);
|
||||
if (exception.response != null && exception.response?.data != null) {
|
||||
_logPrintError('╔ ${exception.type.toString()}');
|
||||
_printResponse(_logPrintError, exception.response!);
|
||||
}
|
||||
_printLine(_logPrintError, '╚');
|
||||
_logPrintError('');
|
||||
} else {
|
||||
_printBoxed(
|
||||
_logPrintError,
|
||||
header: 'DioError ║ ${err.type}',
|
||||
text: err.message,
|
||||
header: 'DioException ║ ${exception.type}',
|
||||
text: exception.message,
|
||||
);
|
||||
_printRequestHeader(_logPrintError, err.requestOptions);
|
||||
_printRequestHeader(_logPrintError, exception.requestOptions);
|
||||
}
|
||||
}
|
||||
super.onError(err, handler);
|
||||
super.onError(exception, handler);
|
||||
}
|
||||
|
||||
@override
|
||||
|
||||
@@ -2,7 +2,7 @@ import 'package:dio/dio.dart';
|
||||
import 'package:stream_chat/src/core/error/error.dart';
|
||||
|
||||
/// Error class specific to StreamChat and Dio
|
||||
class StreamChatDioError extends DioError {
|
||||
class StreamChatDioError extends DioException {
|
||||
/// Initialize a stream chat dio error
|
||||
StreamChatDioError({
|
||||
required this.error,
|
||||
|
||||
@@ -92,16 +92,16 @@ class StreamHttpClient {
|
||||
/// calling [close] will throw an exception.
|
||||
void close({bool force = false}) => httpClient.close(force: force);
|
||||
|
||||
StreamChatNetworkError _parseError(DioError err) {
|
||||
StreamChatNetworkError _parseError(DioException exception) {
|
||||
StreamChatNetworkError error;
|
||||
// locally thrown dio error
|
||||
if (err is StreamChatDioError) {
|
||||
error = err.error;
|
||||
if (exception is StreamChatDioError) {
|
||||
error = exception.error;
|
||||
} else {
|
||||
// real network request dio error
|
||||
error = StreamChatNetworkError.fromDioError(err);
|
||||
error = StreamChatNetworkError.fromDioException(exception);
|
||||
}
|
||||
return error..stackTrace = err.stackTrace;
|
||||
return error..stackTrace = exception.stackTrace;
|
||||
}
|
||||
|
||||
/// Handy method to make http GET request with error parsing.
|
||||
@@ -121,7 +121,7 @@ class StreamHttpClient {
|
||||
cancelToken: cancelToken,
|
||||
);
|
||||
return response;
|
||||
} on DioError catch (error) {
|
||||
} on DioException catch (error) {
|
||||
throw _parseError(error);
|
||||
}
|
||||
}
|
||||
@@ -147,7 +147,7 @@ class StreamHttpClient {
|
||||
cancelToken: cancelToken,
|
||||
);
|
||||
return response;
|
||||
} on DioError catch (error) {
|
||||
} on DioException catch (error) {
|
||||
throw _parseError(error);
|
||||
}
|
||||
}
|
||||
@@ -167,7 +167,7 @@ class StreamHttpClient {
|
||||
cancelToken: cancelToken,
|
||||
);
|
||||
return response;
|
||||
} on DioError catch (error) {
|
||||
} on DioException catch (error) {
|
||||
throw _parseError(error);
|
||||
}
|
||||
}
|
||||
@@ -193,7 +193,7 @@ class StreamHttpClient {
|
||||
cancelToken: cancelToken,
|
||||
);
|
||||
return response;
|
||||
} on DioError catch (error) {
|
||||
} on DioException catch (error) {
|
||||
throw _parseError(error);
|
||||
}
|
||||
}
|
||||
@@ -219,7 +219,7 @@ class StreamHttpClient {
|
||||
cancelToken: cancelToken,
|
||||
);
|
||||
return response;
|
||||
} on DioError catch (error) {
|
||||
} on DioException catch (error) {
|
||||
throw _parseError(error);
|
||||
}
|
||||
}
|
||||
@@ -268,7 +268,7 @@ class StreamHttpClient {
|
||||
cancelToken: cancelToken,
|
||||
);
|
||||
return response;
|
||||
} on DioError catch (error) {
|
||||
} on DioException catch (error) {
|
||||
throw _parseError(error);
|
||||
}
|
||||
}
|
||||
@@ -281,7 +281,7 @@ class StreamHttpClient {
|
||||
try {
|
||||
final response = await httpClient.fetch<T>(requestOptions);
|
||||
return response;
|
||||
} on DioError catch (error) {
|
||||
} on DioException catch (error) {
|
||||
throw _parseError(error);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,6 +17,11 @@ abstract class ChatPersistenceClient {
|
||||
/// Whether the connection is established.
|
||||
bool get isConnected;
|
||||
|
||||
/// The current user id to which the client is connected.
|
||||
///
|
||||
/// Returns `null` if the client is not connected.
|
||||
String? get userId;
|
||||
|
||||
/// Creates a new connection to the client
|
||||
Future<void> connect(String userId);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user