runTask<T> method

Future<T?> runTask<T>(
  1. FutureOr<T> task(
    1. LevitTaskContext context
    ), {
  2. String? id,
  3. TaskPriority priority = TaskPriority.normal,
  4. TaskConflictPolicy conflictPolicy = TaskConflictPolicy.reject,
  5. int retries = 0,
  6. Duration? retryDelay,
  7. bool useExponentialBackoff = true,
  8. double weight = 1.0,
  9. void onError(
    1. Object error,
    2. StackTrace stackTrace
    )?,
  10. TaskCachePolicy<T>? cachePolicy,
  11. LevitTaskMetadata? metadata,
  12. 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;
}