fix(llc): fix retry queue mechanism
This commit is contained in:
@@ -71,14 +71,15 @@ class RetryQueue {
|
|||||||
/// Add a list of messages
|
/// Add a list of messages
|
||||||
void add(List<Message> messages) {
|
void add(List<Message> messages) {
|
||||||
if (messages.isEmpty) return;
|
if (messages.isEmpty) return;
|
||||||
if (_messageQueue.containsAllMessage(messages)) return;
|
if (!_messageQueue.containsAllMessage(messages)) {
|
||||||
|
logger?.info('Adding ${messages.length} messages');
|
||||||
|
final messageList = _messageQueue.toList();
|
||||||
|
// we should not add message if already available in the queue
|
||||||
|
_messageQueue.addAll(messages.where(
|
||||||
|
(it) => !messageList.any((m) => m.id == it.id),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
logger?.info('Adding ${messages.length} messages');
|
|
||||||
final messageList = _messageQueue.toList();
|
|
||||||
// we should not add message if already available in the queue
|
|
||||||
_messageQueue.addAll(messages.where(
|
|
||||||
(it) => !messageList.any((m) => m.id == it.id),
|
|
||||||
));
|
|
||||||
_startRetrying();
|
_startRetrying();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -87,9 +88,9 @@ class RetryQueue {
|
|||||||
_isRetrying = true;
|
_isRetrying = true;
|
||||||
|
|
||||||
logger?.info('Started retrying failed messages');
|
logger?.info('Started retrying failed messages');
|
||||||
while (_messageQueue.isNotEmpty) {
|
for (var i = 0; i < _messageQueue.length; ++i) {
|
||||||
logger?.info('${_messageQueue.length} messages remaining in the queue');
|
logger?.info('${_messageQueue.length} messages remaining in the queue');
|
||||||
final message = _messageQueue.first;
|
final message = _messageQueue.toList()[i];
|
||||||
await _runAndRetry(message);
|
await _runAndRetry(message);
|
||||||
}
|
}
|
||||||
_isRetrying = false;
|
_isRetrying = false;
|
||||||
@@ -109,8 +110,12 @@ class RetryQueue {
|
|||||||
await _retryMessage(message);
|
await _retryMessage(message);
|
||||||
logger?.info('Message (${message.id}) sent successfully');
|
logger?.info('Message (${message.id}) sent successfully');
|
||||||
_messageQueue.removeMessage(message);
|
_messageQueue.removeMessage(message);
|
||||||
break;
|
return;
|
||||||
} on StreamChatError catch (e) {
|
} catch (e) {
|
||||||
|
if (e is! StreamChatNetworkError || !e.isRetriable) {
|
||||||
|
_messageQueue.removeMessage(message);
|
||||||
|
return;
|
||||||
|
}
|
||||||
// retry logic
|
// retry logic
|
||||||
final maxAttempt = _retryPolicy.maxRetryAttempts;
|
final maxAttempt = _retryPolicy.maxRetryAttempts;
|
||||||
if (attempt < maxAttempt) {
|
if (attempt < maxAttempt) {
|
||||||
@@ -143,14 +148,6 @@ class RetryQueue {
|
|||||||
_sendFailedEvent(message);
|
_sendFailedEvent(message);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
} catch (e) {
|
|
||||||
logger?.info(
|
|
||||||
'API call failed due to unknown error (attempt $attempt). '
|
|
||||||
'Giving up for now, will retry when connection recovers. '
|
|
||||||
'Error was $e',
|
|
||||||
);
|
|
||||||
_sendFailedEvent(message);
|
|
||||||
break;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -832,8 +832,8 @@ class _MessageWidgetState extends State<MessageWidget>
|
|||||||
),
|
),
|
||||||
if (isFailedState)
|
if (isFailedState)
|
||||||
Positioned(
|
Positioned(
|
||||||
left: widget.reverse ? 0 : null,
|
right: widget.reverse ? 0 : null,
|
||||||
right: widget.reverse ? null : 0,
|
left: widget.reverse ? null : 0,
|
||||||
bottom: showBottomRow ? 18 : -2,
|
bottom: showBottomRow ? 18 : -2,
|
||||||
child: StreamSvgIcon.error(size: 20),
|
child: StreamSvgIcon.error(size: 20),
|
||||||
),
|
),
|
||||||
|
|||||||
Reference in New Issue
Block a user