loadOutputResources<T> function

Future<void> loadOutputResources<T>(
  1. Iterable<T> resources,
  2. Future<Object?> load(
    1. T
    ), {
  3. required int concurrency,
  4. required Duration timeout,
})

Bounded resource admission with one deadline for the whole operation. Active loads are observed after failure, but no queued load starts afterward. The caller owns cancellation/disposal of active resources.

Implementation

Future<void> loadOutputResources<T>(
  Iterable<T> resources,
  Future<Object?> Function(T) load, {
  required int concurrency,
  required Duration timeout,
}) async {
  if (concurrency <= 0 || timeout <= Duration.zero) {
    throw ArgumentError('Resource concurrency and timeout must be positive.');
  }
  final iterator = resources.iterator;
  final done = Completer<void>();
  var stopped = false;
  var workers = 0;
  void fail(Object error, StackTrace stack) {
    if (stopped) return;
    stopped = true;
    done.completeError(error, stack);
  }

  final timer = Timer(timeout, () {
    fail(
      TimeoutException('Output resources timed out.', timeout),
      StackTrace.current,
    );
  });
  Future<void> work(T first) async {
    try {
      await load(first);
      while (!stopped && iterator.moveNext()) {
        await load(iterator.current);
      }
    } catch (error, stack) {
      fail(error, stack);
    } finally {
      workers--;
      if (workers == 0 && !stopped) {
        stopped = true;
        done.complete();
      }
    }
  }

  try {
    for (var i = 0; i < concurrency && !stopped && iterator.moveNext(); i++) {
      workers++;
      unawaited(work(iterator.current));
    }
    if (workers == 0 && !stopped) {
      stopped = true;
      done.complete();
    }
  } catch (error, stack) {
    fail(error, stack);
  }
  try {
    await done.future;
  } finally {
    timer.cancel();
  }
}