fix persistence abstract class

This commit is contained in:
Salvatore Giordano
2021-04-19 16:21:32 +02:00
parent b69be11e4d
commit b250bdb98e
3 changed files with 36 additions and 38 deletions
+30 -30
View File
@@ -28,32 +28,32 @@ class WebSocket {
/// To connect the WS call [connect] /// To connect the WS call [connect]
WebSocket({ WebSocket({
required this.baseUrl, required this.baseUrl,
this.user, required this.user,
this.connectParams, required this.handler,
this.connectPayload, this.connectParams = const {},
this.handler, this.connectPayload = const {},
this.logger, this.logger,
this.connectFunc, this.connectFunc,
this.reconnectionMonitorInterval = 1, this.reconnectionMonitorInterval = 1,
this.healthCheckInterval = 20, this.healthCheckInterval = 20,
this.reconnectionMonitorTimeout = 40, this.reconnectionMonitorTimeout = 40,
}) { }) {
final qs = Map<String, String>.from(connectParams!); final qs = Map<String, String>.from(connectParams);
final data = Map<String, dynamic>.from(connectPayload!); final data = Map<String, dynamic>.from(connectPayload);
data['user_details'] = user!.toJson(); data['user_details'] = user.toJson();
qs['json'] = json.encode(data); qs['json'] = json.encode(data);
if (baseUrl.startsWith('https')) { if (baseUrl.startsWith('https')) {
_path = baseUrl.replaceFirst('https://', ''); _path = baseUrl.replaceFirst('https://', '');
_path = Uri.https(_path!, 'connect', qs) _path = Uri.https(_path, 'connect', qs)
.toString() .toString()
.replaceFirst('https', 'wss'); .replaceFirst('https', 'wss');
} else if (baseUrl.startsWith('http')) { } else if (baseUrl.startsWith('http')) {
_path = baseUrl.replaceFirst('http://', ''); _path = baseUrl.replaceFirst('http://', '');
_path = _path =
Uri.http(_path!, 'connect', qs).toString().replaceFirst('http', 'ws'); Uri.http(_path, 'connect', qs).toString().replaceFirst('http', 'ws');
} else { } else {
_path = Uri.https(baseUrl, 'connect', qs) _path = Uri.https(baseUrl, 'connect', qs)
.toString() .toString()
@@ -65,17 +65,17 @@ class WebSocket {
final String baseUrl; final String baseUrl;
/// User performing the WS connection /// User performing the WS connection
final User? user; final User user;
/// Querystring connection parameters /// Querystring connection parameters
final Map<String, String?>? connectParams; final Map<String, String> connectParams;
/// WS connection payload /// WS connection payload
final Map<String, dynamic>? connectPayload; final Map<String, dynamic> connectPayload;
/// Functions that will be called every time a new event is received from the /// Functions that will be called every time a new event is received from the
/// connection /// connection
final EventHandler? handler; final EventHandler handler;
/// A WS specific logger instance /// A WS specific logger instance
final Logger? logger; final Logger? logger;
@@ -113,7 +113,7 @@ class WebSocket {
Stream<ConnectionStatus> get connectionStatusStream => Stream<ConnectionStatus> get connectionStatusStream =>
_connectionStatusController.stream; _connectionStatusController.stream;
String? _path; late final String _path;
int _retryAttempt = 1; int _retryAttempt = 1;
late WebSocketChannel _channel; late WebSocketChannel _channel;
Timer? _healthCheck, _reconnectionMonitor; Timer? _healthCheck, _reconnectionMonitor;
@@ -131,17 +131,17 @@ class WebSocket {
_manuallyDisconnected = false; _manuallyDisconnected = false;
if (_connecting) { if (_connecting) {
logger!.severe('already connecting'); logger?.severe('already connecting');
return null; return null;
} }
_connecting = true; _connecting = true;
_connectionStatus = ConnectionStatus.connecting; _connectionStatus = ConnectionStatus.connecting;
logger!.info('connecting to $_path'); logger?.info('connecting to $_path');
_channel = _channel =
connectFunc?.call(_path) ?? WebSocketChannel.connect(Uri.parse(_path!)); connectFunc?.call(_path) ?? WebSocketChannel.connect(Uri.parse(_path));
_channel.stream.listen( _channel.stream.listen(
(data) async { (data) async {
final jsonData = json.decode(data); final jsonData = json.decode(data);
@@ -164,7 +164,7 @@ class WebSocket {
return; return;
} }
logger!.info('connection closed | closeCode: ${_channel.closeCode} | ' logger?.info('connection closed | closeCode: ${_channel.closeCode} | '
'closedReason: ${_channel.closeReason}'); 'closedReason: ${_channel.closeReason}');
if (!_reconnecting) { if (!_reconnecting) {
@@ -178,10 +178,10 @@ class WebSocket {
} }
final event = _decodeEvent(data); final event = _decodeEvent(data);
logger!.info('received new event: $data'); logger?.info('received new event: $data');
if (_lastEventAt == null) { if (_lastEventAt == null) {
logger!.info('connection estabilished'); logger?.info('connection estabilished');
_connecting = false; _connecting = false;
_reconnecting = false; _reconnecting = false;
_lastEventAt = DateTime.now(); _lastEventAt = DateTime.now();
@@ -197,14 +197,14 @@ class WebSocket {
_startHealthCheck(); _startHealthCheck();
} }
handler!(event); handler(event);
_lastEventAt = DateTime.now(); _lastEventAt = DateTime.now();
} }
Future<void> _onConnectionError(error, [stacktrace]) async { Future<void> _onConnectionError(error, [stacktrace]) async {
logger!..severe('error connecting')..severe(error); logger?..severe('error connecting')..severe(error);
if (stacktrace != null) { if (stacktrace != null) {
logger!.severe(stacktrace); logger?.severe(stacktrace);
} }
_connecting = false; _connecting = false;
@@ -242,18 +242,18 @@ class WebSocket {
return; return;
} }
if (_connecting) { if (_connecting) {
logger!.info('already connecting'); logger?.info('already connecting');
return; return;
} }
logger!.info('reconnecting..'); logger?.info('reconnecting..');
_cancelTimers(); _cancelTimers();
try { try {
await connect(); await connect();
} catch (e) { } catch (e) {
logger!.log(Level.SEVERE, e.toString()); logger?.log(Level.SEVERE, e.toString());
} }
await Future.delayed( await Future.delayed(
Duration(seconds: min(_retryAttempt * 5, 25)), Duration(seconds: min(_retryAttempt * 5, 25)),
@@ -265,7 +265,7 @@ class WebSocket {
} }
Future<void> _reconnect() async { Future<void> _reconnect() async {
logger!.info('reconnect'); logger?.info('reconnect');
if (!_reconnecting) { if (!_reconnecting) {
_reconnecting = true; _reconnecting = true;
_connectionStatus = ConnectionStatus.connecting; _connectionStatus = ConnectionStatus.connecting;
@@ -285,12 +285,12 @@ class WebSocket {
} }
void _healthCheckTimer(_) { void _healthCheckTimer(_) {
logger!.info('sending health.check'); logger?.info('sending health.check');
_channel.sink.add("{'type': 'health.check'}"); _channel.sink.add("{'type': 'health.check'}");
} }
void _startHealthCheck() { void _startHealthCheck() {
logger!.info('start health check monitor'); logger?.info('start health check monitor');
_healthCheck = Timer.periodic( _healthCheck = Timer.periodic(
Duration(seconds: healthCheckInterval), Duration(seconds: healthCheckInterval),
@@ -309,7 +309,7 @@ class WebSocket {
if (_manuallyDisconnected) { if (_manuallyDisconnected) {
return; return;
} }
logger!.info('disconnecting'); logger?.info('disconnecting');
_connectionCompleter = Completer(); _connectionCompleter = Completer();
_cancelTimers(); _cancelTimers();
_reconnecting = false; _reconnecting = false;
+2 -2
View File
@@ -523,10 +523,10 @@ class StreamChatClient {
_ws = WebSocket( _ws = WebSocket(
baseUrl: baseURL, baseUrl: baseURL,
user: state.user, user: state.user!,
connectParams: { connectParams: {
'api_key': apiKey, 'api_key': apiKey,
'authorization': token, 'authorization': token!,
'stream-auth-type': _authType, 'stream-auth-type': _authType,
'X-Stream-Client': _userAgent, 'X-Stream-Client': _userAgent,
}, },
@@ -199,11 +199,9 @@ abstract class ChatPersistenceClient {
final reactions = final reactions =
cleanedChannelStates.expand((it) => it.messages).expand((it) => [ cleanedChannelStates.expand((it) => it.messages).expand((it) => [
if (it.ownReactions != null) if (it.ownReactions != null)
...it.ownReactions?.where((r) => r.userId != null) ?? ...it.ownReactions!.where((r) => r.userId != null),
<Reaction>[],
if (it.latestReactions != null) if (it.latestReactions != null)
...it.latestReactions?.where((r) => r.userId != null) ?? ...it.latestReactions!.where((r) => r.userId != null),
<Reaction>[],
]); ]);
final users = cleanedChannelStates final users = cleanedChannelStates
@@ -213,9 +211,9 @@ abstract class ChatPersistenceClient {
.map((m) => [ .map((m) => [
m.user, m.user,
if (m.latestReactions != null) if (m.latestReactions != null)
...m.latestReactions?.map((r) => r.user) ?? [], ...m.latestReactions!.map((r) => r.user),
if (m.ownReactions != null) if (m.ownReactions != null)
...m.ownReactions?.map((r) => r.user) ?? [], ...m.ownReactions!.map((r) => r.user),
]) ])
.expand((v) => v), .expand((v) => v),
...cs.read.map((r) => r.user), ...cs.read.map((r) => r.user),