feat: safe cancellation of running operation on reset

This commit is contained in:
Kingkor Roy Tirtho
2023-02-23 13:14:42 +06:00
parent 68e60c1914
commit 5dde200c84
9 changed files with 67 additions and 37 deletions
@@ -1,9 +1,10 @@
import 'dart:async';
import 'package:async/async.dart';
import 'package:collection/collection.dart';
import 'package:fl_query/fl_query.dart';
import 'package:fl_query/src/core/retryer.dart';
import 'package:fl_query/src/core/validation.dart';
import 'package:fl_query/src/core/mixins/retryer.dart';
import 'package:fl_query/src/core/mixins/validation.dart';
import 'package:hive_flutter/adapters.dart';
import 'package:mutex/mutex.dart';
import 'package:state_notifier/state_notifier.dart';
@@ -108,6 +109,8 @@ class InfiniteQuery<DataType, ErrorType, PageType>
final RefreshConfig refreshConfig;
final JsonConfig<DataType>? jsonConfig;
final PageType _initialParam;
InfiniteQuery(
this.key,
InfiniteQueryFn<DataType, PageType> queryFn, {
@@ -116,7 +119,8 @@ class InfiniteQuery<DataType, ErrorType, PageType>
this.retryConfig = DefaultConstants.retryConfig,
this.refreshConfig = DefaultConstants.refreshConfig,
this.jsonConfig,
}) : _dataController = StreamController.broadcast(),
}) : _initialParam = initialParam,
_dataController = StreamController.broadcast(),
_errorController = StreamController.broadcast(),
_box = Hive.lazyBox(QueryClient.infiniteQueryCachePrefix),
super(InfiniteQueryState<DataType, ErrorType, PageType>(
@@ -175,6 +179,8 @@ class InfiniteQuery<DataType, ErrorType, PageType>
final StreamController<PageEvent<DataType, PageType>> _dataController;
final StreamController<PageEvent<ErrorType, PageType>> _errorController;
CancelableOperation<void>? _operation;
List<DataType> get pages =>
state.pages.map((e) => e.data).whereType<DataType>().toList();
List<ErrorType> get errors =>
@@ -197,10 +203,10 @@ class InfiniteQuery<DataType, ErrorType, PageType>
bool get hasNextPage => state.hasNextPage;
Future<void> _operation(PageType page) {
Future<void> _operate(PageType page) {
return _mutex.protect(() async {
state = state.copyWith();
return await retryOperation(
_operation = cancellableRetryOperation(
() => state.queryFn(page),
config: retryConfig,
onSuccessful: (data) async {
@@ -262,14 +268,14 @@ class InfiniteQuery<DataType, ErrorType, PageType>
final lastPage = state.lastPage;
if (_mutex.isLocked || hasPageData || hasPageError)
return state.pages.last.data;
return await _operation(lastPage).then((_) => state.pages.last.data);
return await _operate(lastPage).then((_) => state.pages.last.data);
}
Future<DataType?> refresh([PageType? page]) async {
page ??= lastPage;
if (_mutex.isLocked)
return state.pages.firstWhereOrNull((e) => e.page == page)?.data;
return await _operation(page!).then((_) {
return await _operate(page!).then((_) {
return state.pages.firstWhereOrNull((e) => e.page == page)?.data;
});
}
@@ -277,7 +283,7 @@ class InfiniteQuery<DataType, ErrorType, PageType>
Future<List<DataType>?> refreshAll() async {
if (_mutex.isLocked) return pages;
return await Future.wait(
state.pages.map((e) => _operation(e.page)),
state.pages.map((e) => _operate(e.page)),
).then((_) => pages);
}
@@ -286,7 +292,7 @@ class InfiniteQuery<DataType, ErrorType, PageType>
if (_mutex.isLocked || nextPage == null) {
return state.pages.lastOrNull?.data;
}
return await _operation(nextPage).then((_) {
return await _operate(nextPage).then((_) {
return state.pages.firstWhereOrNull((e) => e.page == nextPage)?.data;
});
}
@@ -331,6 +337,18 @@ class InfiniteQuery<DataType, ErrorType, PageType>
);
}
Future<void> reset() async {
await _operation?.cancel();
state = state.copyWith(pages: {
InfiniteQueryPage<DataType, ErrorType, PageType>(
page: _initialParam,
updatedAt: DateTime.now(),
staleDuration: refreshConfig.staleDuration,
)
});
_box.delete(key);
}
@override
RemoveListener addListener(
Listener<InfiniteQueryState<DataType, ErrorType, PageType>> listener, {
@@ -1,5 +1,6 @@
import 'dart:async';
import 'package:async/async.dart';
import 'package:fl_query/src/collections/retry_config.dart';
import 'package:flutter/material.dart';
@@ -38,4 +39,20 @@ mixin Retryer<T, E> {
}
}
}
CancelableOperation<void> cancellableRetryOperation(
FutureOr<T?> Function() operation, {
required RetryConfig config,
required void Function(T?) onSuccessful,
required void Function(E?) onFailed,
}) {
return CancelableOperation.fromFuture(
retryOperation(
operation,
config: config,
onSuccessful: onSuccessful,
onFailed: onFailed,
),
);
}
}
+6 -3
View File
@@ -1,8 +1,9 @@
import 'dart:async';
import 'package:async/async.dart';
import 'package:fl_query/src/collections/default_configs.dart';
import 'package:fl_query/src/collections/retry_config.dart';
import 'package:fl_query/src/core/retryer.dart';
import 'package:fl_query/src/core/mixins/retryer.dart';
import 'package:mutex/mutex.dart';
import 'package:state_notifier/state_notifier.dart';
@@ -74,11 +75,12 @@ class Mutation<DataType, ErrorType, VariablesType>
final StreamController<VariablesType> _mutationController;
final StreamController<DataType> _dataController;
final StreamController<ErrorType> _errorController;
CancelableOperation<void>? _operation;
Future<void> _operate(VariablesType variables) {
return _mutex.protect(() async {
state = state.copyWith();
return await retryOperation(
_operation = await cancellableRetryOperation(
() {
_mutationController.add(variables);
return state.mutationFn(variables);
@@ -115,7 +117,8 @@ class Mutation<DataType, ErrorType, VariablesType>
state = state.copyWith(mutationFn: mutationFn, updatedAt: state.updatedAt);
}
void reset() {
Future<void> reset() async {
await _operation?.cancel();
state = MutationState<DataType, ErrorType, VariablesType>(
mutationFn: state.mutationFn,
);
+14 -5
View File
@@ -5,11 +5,12 @@ import 'package:fl_query/src/collections/json_config.dart';
import 'package:fl_query/src/collections/refresh_config.dart';
import 'package:fl_query/src/collections/retry_config.dart';
import 'package:fl_query/src/core/client.dart';
import 'package:fl_query/src/core/retryer.dart';
import 'package:fl_query/src/core/validation.dart';
import 'package:fl_query/src/core/mixins/retryer.dart';
import 'package:fl_query/src/core/mixins/validation.dart';
import 'package:hive_flutter/adapters.dart';
import 'package:mutex/mutex.dart';
import 'package:state_notifier/state_notifier.dart';
import 'package:async/async.dart';
typedef QueryFn<DataType> = FutureOr<DataType?> Function();
@@ -73,7 +74,7 @@ class Query<DataType, ErrorType>
)) {
if (jsonConfig != null) {
_mutex.protect(() async {
final json = await _box.get(key.toString());
final json = await _box.get(key);
if (json != null) {
_initial = jsonConfig!.fromJson(
Map.castFrom<dynamic, dynamic, String, dynamic>(json),
@@ -116,10 +117,12 @@ class Query<DataType, ErrorType>
Stream<DataType> get dataStream => _dataController.stream;
Stream<ErrorType> get errorStream => _errorController.stream;
CancelableOperation<void>? _operation;
Future<void> _operate() {
return _mutex.protect(() async {
state = state.copyWith();
return await retryOperation(
_operation = cancellableRetryOperation(
state.queryFn,
config: retryConfig,
onSuccessful: (DataType? data) {
@@ -131,7 +134,7 @@ class Query<DataType, ErrorType>
_dataController.add(data);
if (jsonConfig != null) {
_box.put(
key.toString(),
key,
jsonConfig!.toJson(data),
);
}
@@ -168,6 +171,12 @@ class Query<DataType, ErrorType>
state = state.copyWith(data: data, updatedAt: DateTime.now());
}
Future<void> reset() async {
await _operation?.cancel();
state = state.copyWith(data: _initial, updatedAt: DateTime.now());
_box.delete(key);
}
@override
RemoveListener addListener(Listener<QueryState<DataType, ErrorType>> listener,
{bool fireImmediately = true}) {