Перейти к содержимому

Дети, стримы и цепочки

Тело поручает работу другому Job, зарегистрировав его как дочерний:

final parent = Job<void>((ctx) async {
final child = Job.deferred<int>((ctx) => ctx.wait(load));
final rows = await ctx.run(child);
ctx.log('$rows rows');
});

ctx.run(child) сразу запускает дочернюю задачу и возвращает Future её результата. Родитель ждёт всех дочерних задач перед завершением и передаёт им отмену. Дочерняя задача с cancellable: false может отклонить эту отмену. Если тело родителя бросит Cancelled, дочерние задачи тоже будут отменены. При другой ошибке родитель позволяет им доработать и ждёт их завершения.

Создавайте дочерние задачи через Job.deferred, чтобы родитель управлял их запуском. ctx.run отклоняет обычный Job, даже до запланированного старта. Это исключает гонку, при которой дочерняя задача запускается самостоятельно и остаётся вне отмены и ожидания родителя.

run ждёт завершения дочерней задачи, включая её детей и уборку, и проверяет отмену родителя перед возвратом значения в тело. Если дочерняя задача отклонит отмену, run продолжит ждать её; после её успешного завершения он бросит отмену родителя, и тело не перейдёт к вызову log. Ошибка или отмена дочерней задачи передаётся через возвращённую Future со своим стеком.

Для запроса отмены вызовите child.cancel(), для просмотра исхода используйте child.done. Для параллельного запуска сохраните Future от ctx.run(child) и дождитесь её позже либо обработайте её ошибки. Если результат намеренно не нужен, используйте ctx.run(child).ignore(). Родитель всё равно дождётся ребёнка перед завершением. Одного child.ignore() недостаточно для обработки ошибок Future от run; необработанная ошибка Future, включая отмену, подчиняется обычным правилам Dart.

Дочерняя задача наследует наблюдателя родителя, если у неё нет своего. Если её отмена выходит через await ctx.run(child) или child.value, родитель завершается с HandlerCancelReason и описанием, указывающим на дочернюю задачу.

При недопустимом запуске ctx.run бросает синхронно: ArgumentError для задачи другой реализации или задачи с автоматическим запуском. Он бросает StateError, если дочерняя задача уже запущена или тело родителя завершилось. Если родитель уже отменён, метод отменяет дочернюю задачу до старта и бросает Cancelled родителя.

ctx.each(stream, onData) подписывается на поток и обрабатывает события по одному. Метод сразу возвращает дочерний Job<void>, которому принадлежит подписка. Колбэк получает контекст этого ребёнка и событие.

Например, saveMessages ниже принимает поток сообщений и асинхронную функцию сохранения одного сообщения. each ждёт завершения каждого сохранения, прежде чем передать следующее сообщение в колбэк:

Future<void> saveMessages(
Stream<String> messages,
Future<void> Function(String message) save,
) async {
final job = Job<void>((ctx) async {
final processing = ctx.each(messages, (child, message) {
return child.join(() => save(message));
});
await processing.value;
});
await job.value;
}

processing.value завершается, когда поток закончился и последнее сохранение выполнено. Ошибка потока или колбэка прекращает обработку; ожидание value передаёт её в тело родителя и затем вызывающему saveMessages. Брошенный Cancelled идёт путём отмены. Внутри колбэка используйте контекст ребёнка, как child.join в примере: он проверяет отмену ребёнка до и после операции сохранения.

Возвращённый Job также позволяет родителю отдельно прекратить подписку. Этот пример печатает событие каждую секунду, прекращает подписку через 2,5 секунды и продолжает тело родителя:

Future<void> watchTicks() async {
final job = Job<void>((ctx) async {
final ticks = ctx.each(
Stream<int>.periodic(const Duration(seconds: 1), (index) => index + 1),
(child, tick) => print('Tick $tick'),
);
await ctx.wait(
() => Future<void>.delayed(const Duration(milliseconds: 2500)),
);
await ticks.cancel();
print('Parent continues');
});
await job.value;
}

Вывод: Tick 1, Tick 2, затем Parent continues. Поток ещё не завершён, когда ticks.cancel() снимает подписку. Вызов ждёт завершения ребёнка; отмена этого ребёнка сама не отменяет родителя.

Используйте childJob.value, когда ошибка или отмена ребёнка должна быть брошена в тело, либо childJob.done, чтобы прочитать исход без броска. Неперехваченный Cancelled из value ребёнка отменяет родителя по обычным правилам исходов дочернего Job.

Родитель ждёт ребёнка, даже если его тело вернулось, не ожидая ребёнка. Поэтому открытый поток удерживает родителя до конца потока или отмены ребёнка. Принятая отмена родителя передаётся ребёнку. Наблюдатель видит этого ребёнка как отдельный Job.

Приняв отмену, ребёнок сразу снимает подписку и прекращает доставку событий. Затем он дожидается уже начатого колбэка, прежде чем завершиться и освободить ресурсы; уборка родителя тоже ждёт. Используйте внутри колбэка контрольные точки дочернего контекста: голый await нельзя прервать, и он может задержать отмену навсегда. Не ждите завершения самого ребёнка или его cancel() из колбэка: ребёнок уже ждёт этот колбэк.

Future, которую возвращает cancel() самой подписки, не ожидается. Если источнику нужна асинхронная уборка, организуйте её ожидание отдельно. Обычное завершение потока по-прежнему зависит от отправки onDone источником. Как и run, each запускает детей только пока тело родителя активно; вызовы из unattended или уборки отвергаются.

then продолжает работу с результатом предыдущего Job. Каждый вызов возвращает новый Job и передаёт обработчику отдельный контекст. Следующий обработчик запускается после успешного завершения предыдущего Job, включая его детей и уборку:

final loaded = Job<String>((ctx) => ctx.join(loadText));
final parsed = loaded.then<int>((ctx, text) => int.parse(text));
final saved = parsed.then<void>((ctx, number) => ctx.join(
() => saveNumber(number),
));
await saved.value;

Отмена saved запрашивает отмену ещё не завершённых parsed и loaded. Отмена parsed запрашивает отмену незавершённого loaded и отменяет saved. Отмена loaded передаётся вперёд обоим продолжениям. Переданный запрос получает ChainCancelReason, в поле cause которой лежит Cancelled соседнего звена. Завершённые Job сохраняют свой исход. Если к одному Job присоединено несколько продолжений, отмена одного может отменить общий источник и остальные продолжения.

await saved.cancel() ждёт незавершённых предшественников, их детей и уборку, а также уже начатую работу и уборку самого saved. Источник сохраняет свои правила отмены: cancellable: false отклоняет запрос, а uncancellable придерживает его. Отменённое продолжение всё равно ждёт такой источник, затем завершается без вызова обработчика. Во время ожидания у него может быть isCancelled, но ещё не isFinished; исход содержит started: false, если отмена принята при ожидании источника.

Провал предшественника передаёт ошибку и стек дальше без вызова обработчика. Наблюдайте последнее звено через value, done или ignore, чтобы обработать этот провал. Если отменённое продолжение не может передать ошибку источника, оно её и не наблюдает: например, источник, который отказал отмене, а затем упал, требует отдельной обработки ошибки.

Обработчик возвращает значение или future. Возврат другого Job не ожидает его: для отложенного ребёнка верните ctx.run(child). then не запускает отложенный источник, а само продолжение нельзя усыновить через ctx.run. Не ожидайте продолжение из тела или уборки его источника: продолжение само ждёт завершения этого источника.

Каждое продолжение — корневой Job ядра. У него есть необязательный аргумент observer; наблюдатель, состояние, правила и слот очереди источника не наследуются. Зарегистрированная источником уборка уже выполнена к передаче значения: ресурс, закрытый через onDispose источника, в этот момент уже закрыт.