runTask<T> method
Future<T?>
runTask<T>(
- FutureOr<
T> task(- LevitTaskContext context
- String? id,
- TaskPriority priority = TaskPriority.normal,
- TaskConflictPolicy conflictPolicy = TaskConflictPolicy.reject,
- int retries = 0,
- Duration? retryDelay,
- bool useExponentialBackoff = true,
- double weight = 1.0,
- void onError(
- Object error,
- StackTrace stackTrace
- TaskCachePolicy<
T> ? cachePolicy, - LevitTaskMetadata? metadata,
- String? debugName,
Executes task and tracks the latest execution for its logical ID.
Implementation
Future<T?> runTask<T>(
FutureOr<T> Function(LevitTaskContext context) task, {
String? id,
TaskPriority priority = TaskPriority.normal,
TaskConflictPolicy conflictPolicy = TaskConflictPolicy.reject,
int retries = 0,
Duration? retryDelay,
bool useExponentialBackoff = true,
double weight = 1.0,
void Function(Object error, StackTrace stackTrace)? onError,
TaskCachePolicy<T>? cachePolicy,
LevitTaskMetadata? metadata,
String? debugName,
}) {
_ensureReactiveTaskState();
if (!weight.isFinite || weight < 0) {
throw RangeError.value(weight, 'weight');
}
final taskId = id ?? LevitTaskEngine._generateTaskId();
final taskMetadata = metadata ??
(debugName == null
? LevitTaskMetadata.none
: LevitTaskMetadata(debugName: debugName));
LevitTaskEvent? queuedEvent;
late final LevitTaskExecution<T> execution;
execution = tasksEngine.submit<T>(
task,
id: taskId,
priority: priority,
conflictPolicy: conflictPolicy,
retries: retries,
retryDelay: retryDelay,
useExponentialBackoff: useExponentialBackoff,
cachePolicy: cachePolicy,
metadata: taskMetadata,
onEvent: (event) {
queuedEvent ??= event;
_applyTaskEvent(event);
},
onSuccess: (result) {
_updateExecution(taskId, execution.executionId, (current) {
return current.copyWith(
status: LxSuccess<T>(result),
progress: 1,
);
});
},
onProgress: (progress) {
_updateExecution(taskId, execution.executionId, (current) {
return current.copyWith(progress: progress);
});
},
onError: (error, stackTrace) {
_updateExecution(taskId, execution.executionId, (current) {
return current.copyWith(
status: LxError<Object>(
error,
stackTrace,
current.status.lastValue,
),
);
});
(onError ?? this.onTaskError)?.call(error, stackTrace);
},
onCancel: () {
_updateExecution(taskId, execution.executionId, (current) {
return current.copyWith(
status: LxIdle<dynamic>(current.status.lastValue),
);
});
},
);
if (execution.disposition == LevitTaskSubmissionDisposition.joined ||
execution.disposition == LevitTaskSubmissionDisposition.dropped) {
return execution.result;
}
_cleanupTimers.remove(taskId)?.cancel();
_pruneTaskHistoryFor(taskId);
final initialEvent = queuedEvent;
tasks[taskId] = TaskDetails(
status: LxWaiting<dynamic>(tasks[taskId]?.status.lastValue),
executionId: execution.executionId,
ownerPath: initialEvent?.ownerPath ?? ownerPath,
phase: initialEvent?.phase ?? LevitTaskPhase.queued,
metadata: taskMetadata,
priority: priority,
attempt: initialEvent?.attempt ?? 0,
weight: weight,
progress: 0,
started: false,
queuedAt: initialEvent?.queuedAt ?? DateTime.now(),
);
return execution.result;
}