TF-1821 Add Task and WorkerQueue object

(cherry picked from commit 42a90ce7c6bf773416db9a215e0a4c8b348c85ff)
This commit is contained in:
dab246
2023-05-08 18:14:51 +07:00
committed by Dat Vu
parent d9aa9c4de2
commit ba4d53a833
3 changed files with 131 additions and 0 deletions
@@ -0,0 +1,44 @@
import 'dart:async';
import 'package:equatable/equatable.dart';
import 'package:tmail_ui_user/features/offline_mode/worker/hive_task_state.dart';
class HiveTask<A> with EquatableMixin {
final String? id;
final A? action;
final Future<bool> Function()? conditionInvoked;
final Future Function() runnable;
HiveTask({
required this.runnable,
this.id,
this.action,
this.conditionInvoked,
});
Future execute() async {
final resultCompleter = Completer();
try {
if (conditionInvoked != null) {
final invoked = await conditionInvoked!.call();
if (invoked) {
final result = await runnable.call();
resultCompleter.complete(TaskSuccess(result: result));
} else {
resultCompleter.completeError(TaskFailure());
}
} else {
final result = await runnable.call();
resultCompleter.complete(TaskSuccess(result: result));
}
} catch (e) {
resultCompleter.completeError(TaskFailure(exception: e));
}
return resultCompleter.future;
}
@override
List<Object?> get props => [runnable, id, action, conditionInvoked];
}
@@ -0,0 +1,22 @@
import 'package:core/presentation/state/failure.dart';
import 'package:core/presentation/state/success.dart';
class TaskSuccess extends Success {
final dynamic result;
TaskSuccess({this.result});
@override
List<Object?> get props => [result];
}
class TaskFailure extends Failure {
final dynamic exception;
TaskFailure({this.exception});
@override
List<Object?> get props => [exception];
}
@@ -0,0 +1,65 @@
import 'dart:async';
import 'dart:collection';
import 'package:core/utils/app_logger.dart';
import 'package:tmail_ui_user/features/offline_mode/worker/hive_task.dart';
abstract class WorkerQueue<A> {
final Queue<HiveTask<A>> queue = Queue<HiveTask<A>>();
Completer? completer;
String get workerName;
Future addTask(HiveTask<A> task) {
queue.add(task);
log('WorkerQueue<$workerName>::addTask(): QUEUE_LENGTH: ${queue.length}');
return _processTask();
}
Future _processTask() async {
if (completer != null) {
return completer!.future;
}
completer = Completer();
if (queue.isNotEmpty) {
final firstTask = queue.removeFirst();
log('WorkerQueue<$workerName>::_processTask(): ${firstTask.id}');
firstTask.execute()
.then(_handleTaskExecuteCompleted)
.catchError(_handleTaskExecuteError);
} else {
completer?.complete();
}
return completer!.future;
}
void _handleTaskExecuteCompleted(dynamic value) {
log('WorkerQueue<$workerName>::_handleTaskExecuteCompleted(): $value');
completer?.complete();
_releaseCompleter();
if (queue.isNotEmpty) {
_processTask();
}
}
void _handleTaskExecuteError(error) {
log('WorkerQueue<$workerName>::_handleTaskExecuteError(): $error');
completer?.complete();
_releaseCompleter();
if (queue.isNotEmpty) {
_processTask();
}
}
void _releaseCompleter() {
completer = null;
}
Future release() async {
log('WorkerQueue<$workerName>::release():');
queue.clear();
_releaseCompleter();
}
}