fix: Review changes
This commit is contained in:
@@ -249,13 +249,13 @@ class Channel {
|
||||
Future<String> future;
|
||||
if (isImage) {
|
||||
future = sendImage(
|
||||
it.file,
|
||||
it.file!,
|
||||
onSendProgress: onSendProgress,
|
||||
cancelToken: cancelToken,
|
||||
).then((it) => it!.file!);
|
||||
} else {
|
||||
future = sendFile(
|
||||
it.file,
|
||||
it.file!,
|
||||
onSendProgress: onSendProgress,
|
||||
cancelToken: cancelToken,
|
||||
).then((it) => it!.file!);
|
||||
@@ -340,7 +340,7 @@ class Channel {
|
||||
}
|
||||
|
||||
final response = await _client.sendMessage(message, id, type);
|
||||
state?.addMessage(response!.message!);
|
||||
state?.addMessage(response.message!);
|
||||
return response;
|
||||
} catch (error) {
|
||||
if (error is DioError && error.type != DioErrorType.response) {
|
||||
@@ -392,7 +392,7 @@ class Channel {
|
||||
|
||||
final response = await _client.updateMessage(message);
|
||||
|
||||
final m = response?.message?.copyWith(
|
||||
final m = response.message?.copyWith(
|
||||
ownReactions: message.ownReactions,
|
||||
);
|
||||
|
||||
@@ -453,10 +453,11 @@ class Channel {
|
||||
/// Pins provided message
|
||||
Future<UpdateMessageResponse?> pinMessage(
|
||||
Message message,
|
||||
Object timeoutOrExpirationDate,
|
||||
Object? timeoutOrExpirationDate,
|
||||
) {
|
||||
assert(() {
|
||||
if (timeoutOrExpirationDate is! DateTime &&
|
||||
timeoutOrExpirationDate != null &&
|
||||
timeoutOrExpirationDate is! num) {
|
||||
throw ArgumentError('Invalid timeout or Expiration date');
|
||||
}
|
||||
@@ -485,7 +486,7 @@ class Channel {
|
||||
|
||||
/// Send a file to this channel
|
||||
Future<SendFileResponse?> sendFile(
|
||||
AttachmentFile? file, {
|
||||
AttachmentFile file, {
|
||||
ProgressCallback? onSendProgress,
|
||||
CancelToken? cancelToken,
|
||||
}) =>
|
||||
@@ -499,7 +500,7 @@ class Channel {
|
||||
|
||||
/// Send an image to this channel
|
||||
Future<SendImageResponse?> sendImage(
|
||||
AttachmentFile? file, {
|
||||
AttachmentFile file, {
|
||||
ProgressCallback? onSendProgress,
|
||||
CancelToken? cancelToken,
|
||||
}) =>
|
||||
@@ -765,7 +766,7 @@ class Channel {
|
||||
'message_id': messageId,
|
||||
});
|
||||
|
||||
final res = _client.decode(response.data, SendActionResponse.fromJson)!;
|
||||
final res = _client.decode(response.data, SendActionResponse.fromJson);
|
||||
|
||||
if (res.message != null) {
|
||||
state!.addMessage(res.message!);
|
||||
@@ -880,7 +881,7 @@ class Channel {
|
||||
final repliesResponse = _client.decode<QueryRepliesResponse>(
|
||||
response.data,
|
||||
QueryRepliesResponse.fromJson,
|
||||
)!;
|
||||
);
|
||||
|
||||
state?.updateThreadInfo(parentId, repliesResponse.messages);
|
||||
|
||||
@@ -911,7 +912,7 @@ class Channel {
|
||||
final res = _client.decode<GetMessagesByIdResponse>(
|
||||
response.data,
|
||||
GetMessagesByIdResponse.fromJson,
|
||||
)!;
|
||||
);
|
||||
|
||||
final messages = res.messages;
|
||||
|
||||
@@ -999,8 +1000,7 @@ class Channel {
|
||||
|
||||
try {
|
||||
final response = await _client.post(path, data: payload);
|
||||
final updatedState =
|
||||
_client.decode(response.data, ChannelState.fromJson)!;
|
||||
final updatedState = _client.decode(response.data, ChannelState.fromJson);
|
||||
|
||||
if (_id == null) {
|
||||
_id = updatedState.channel!.id;
|
||||
@@ -1191,7 +1191,7 @@ class Channel {
|
||||
|
||||
/// Call this method to dispose the channel client
|
||||
void dispose() {
|
||||
state!.dispose();
|
||||
state?.dispose();
|
||||
}
|
||||
|
||||
void _checkInitialized() {
|
||||
|
||||
@@ -40,16 +40,17 @@ class RetryQueue {
|
||||
}));
|
||||
}
|
||||
|
||||
final HeapPriorityQueue<Message?> _messageQueue = HeapPriorityQueue(_byDate);
|
||||
final HeapPriorityQueue<Message> _messageQueue = HeapPriorityQueue(_byDate);
|
||||
bool _isRetrying = false;
|
||||
RetryPolicy? _retryPolicy;
|
||||
|
||||
/// Add a list of messages
|
||||
void add(List<Message?> messages) {
|
||||
void add(List<Message> messages) {
|
||||
logger?.info('added ${messages.length} messages');
|
||||
final messageList = _messageQueue.toList();
|
||||
|
||||
_messageQueue.addAll(messages
|
||||
.where((element) => !messageList.any((m) => m!.id == element!.id)));
|
||||
.where((element) => !messageList.any((m) => m.id == element.id)));
|
||||
|
||||
if (_messageQueue.isNotEmpty && !_isRetrying) {
|
||||
_startRetrying();
|
||||
@@ -62,7 +63,7 @@ class RetryQueue {
|
||||
final retryPolicy = _retryPolicy!.copyWith(attempt: 0);
|
||||
|
||||
while (_messageQueue.isNotEmpty) {
|
||||
final message = _messageQueue.first!;
|
||||
final message = _messageQueue.first;
|
||||
try {
|
||||
logger?.info('retry attempt ${retryPolicy.attempt}');
|
||||
await _sendMessage(message);
|
||||
@@ -141,7 +142,7 @@ class RetryQueue {
|
||||
final messageList = _messageQueue.toList();
|
||||
if (event.message != null) {
|
||||
final messageIndex =
|
||||
messageList.indexWhere((m) => m!.id == event.message!.id);
|
||||
messageList.indexWhere((m) => m.id == event.message!.id);
|
||||
if (messageIndex == -1 &&
|
||||
[
|
||||
MessageSendingStatus.failed_update,
|
||||
@@ -149,7 +150,11 @@ class RetryQueue {
|
||||
MessageSendingStatus.failed_delete,
|
||||
].contains(event.message!.status)) {
|
||||
logger?.info('add message from events');
|
||||
add([event.message]);
|
||||
final m = event.message;
|
||||
|
||||
if (m != null) {
|
||||
add([m]);
|
||||
}
|
||||
} else if (messageIndex != -1 &&
|
||||
[
|
||||
MessageSendingStatus.sent,
|
||||
@@ -167,9 +172,13 @@ class RetryQueue {
|
||||
_subscriptions.forEach((s) => s.cancel());
|
||||
}
|
||||
|
||||
static int _byDate(Message? m1, Message? m2) {
|
||||
final date1 = _getMessageDate(m1!)!;
|
||||
final date2 = _getMessageDate(m2!)!;
|
||||
static int _byDate(Message m1, Message m2) {
|
||||
final date1 = _getMessageDate(m1);
|
||||
final date2 = _getMessageDate(m2);
|
||||
|
||||
if (date1 == null || date2 == null) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
return date1.compareTo(date2);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user