classifyRetry method

TaskProcessOutcome classifyRetry(
  1. Envelope envelope,
  2. TaskHandler<Object?> handler,
  3. TaskRetryRequest request,
  4. StackTrace stackTrace,
)

Classifies an explicit retry without publishing or persisting it.

Implementation

TaskProcessOutcome classifyRetry(
  Envelope envelope,
  TaskHandler<Object?> handler,
  TaskRetryRequest request,
  StackTrace stackTrace,
) {
  final policy =
      request.retryPolicy ?? _resolveRetryPolicy(envelope, handler.options);
  final maxRetries =
      request.maxRetries ?? policy?.maxRetries ?? envelope.maxRetries;
  if (envelope.attempt >= maxRetries) {
    return TaskProcessFailure(
      envelope: envelope,
      error: StateError('retry requested but max retries exceeded'),
      stackTrace: stackTrace,
      retryExhausted: true,
    );
  }

  final now = stemNow();
  final requestedAt =
      request.eta ??
      (request.countdown == null ? null : now.add(request.countdown!));
  final computedDelay = requestedAt == null
      ? _computeRetryDelay(envelope.attempt, request, stackTrace, policy)
      : requestedAt.difference(now);
  final delay = computedDelay.isNegative ? Duration.zero : computedDelay;
  final notBefore = requestedAt ?? now.add(delay);
  final updatedMeta = Map<String, Object?>.from(envelope.meta);
  if (request.timeLimit != null) {
    updatedMeta['stem.timeLimitMs'] = request.timeLimit!.inMilliseconds;
  }
  if (request.softTimeLimit != null) {
    updatedMeta['stem.softTimeLimitMs'] =
        request.softTimeLimit!.inMilliseconds;
  }
  if (request.retryPolicy != null) {
    updatedMeta['stem.retryPolicy'] = request.retryPolicy!.toJson();
  }
  final nextEnvelope = envelope.copyWith(
    attempt: envelope.attempt + 1,
    maxRetries: maxRetries,
    notBefore: notBefore,
    meta: updatedMeta,
  );
  return TaskProcessRetry(
    envelope: envelope,
    nextEnvelope: nextEnvelope,
    delay: delay,
    error: request,
    stackTrace: stackTrace,
    explicit: true,
  );
}