run method
Runs until until completes, finishing the job in hand first, and
returns how many jobs completed.
With maxJobs it also returns once that many jobs have completed, or
after a pass over every queue completes none -- dartvel queue work --max-jobs draining a bounded number and going back to the shell. A job
that fails is not a completed one, so a pass of only failures ends it
rather than retrying the same poison job until the bound is reached.
Implementation
Future<int> run({Future<void>? until, int? maxJobs}) async {
if (maxJobs != null && maxJobs < 1) {
throw ArgumentError.value(maxJobs, 'maxJobs', 'must be positive');
}
const DVQueues queue = DVQueues();
if (!queue.adapterConfigured) {
throw StateError(
'This worker has no queue adapter configured. It would work the '
'process-local queue, which no other process can dispatch to, and '
'never receive a job. The generated backend puts the queue on the '
'database DATABASE_URL names; a process started some other way '
'configures one with DVQueues().useAdapter before the worker starts.',
);
}
if (!queue.hasHandlers) {
throw StateError(
'No @DVJob.handler is registered in this worker, so every job it '
'reserved would fail and be dead-lettered. Register the handlers '
'(registerDartvelJobs) before the worker starts.',
);
}
bool stopped = false;
final Completer<void> wake = Completer<void>();
unawaited(
until?.then((_) {
stopped = true;
if (!wake.isCompleted) wake.complete();
}),
);
int total = 0;
while (!stopped) {
int done = 0;
for (final String name in queues) {
if (stopped) break;
final int take =
maxJobs == null ? batch : math.min(batch, maxJobs - total - done);
if (take < 1) break;
done += await queue.work(queue: name, maxJobs: take);
}
total += done;
if (maxJobs != null && (done == 0 || total >= maxJobs)) break;
if (stopped) break;
if (done == 0) {
await Future.any(<Future<void>>[
Future<void>.delayed(idle),
wake.future,
]);
}
}
return total;
}