Аккумуляция событий до запуска задачи
Несколько событий могут разделять одну задачу в очереди. collect
сохраняет каждое событие в списке, accumulate объединяет его с одним
накопленным значением. Обе фабрики возвращают SoloAccumulator<E, T>.
Создайте его один раз в контроллере, затем вызывайте add(event) из
публичных методов контроллера.
Накопитель владеет входными данными ожидающей задачи. Добавление события
не запускает обработчик и не меняет состояние контроллера. Группа принимает
события до закрытия состава: по истечении debounce или при извлечении
задачи из очереди с остальными настройками времени. Обработчик получает
этот вход и обычный SoloContext<S, W>; состояние меняет через ctx.emit.
Аккумуляторы и политики
Заголовок раздела «Аккумуляторы и политики»// По списку на группу: каждое принятое событие сохраняется, по порядку.late final _logs = collect<Ready, LogEntry, void>( (ctx, events) => ctx.join(() => sink.write(events)),);
// По значению на группу: что выживет, решает merge.late final _patches = accumulate<Ready, Patch, void>( (ctx, patch) => ctx.join(() => store.apply(patch)), merge: (previous, incoming) => previous.merge(incoming), // Искать открытую группу этого накопителя где угодно в очереди // и сохранять её место, а не смотреть только на хвост. policy: AccumulationPolicy.join,);
SoloJob<void> log(LogEntry entry) => _logs.add(entry);accumulate объединяет события синхронной функцией merge, которая
определяет, что сохранить; оставить только последнее значение — один
из вариантов. Состояние контроллера обновляет только обработчик Job
из очереди. Работающая группа не принимает новые данные, а группы могут
объединяться только в пределах одного накопителя.
| Политика накопления | Как входные данные присоединяются к ожидающей работе |
|---|---|
AccumulationPolicy.adjacent |
Присоединить к совместимой группе только в хвосте очереди. |
AccumulationPolicy.join |
Присоединить к существующей группе на её текущем месте в очереди. |
AccumulationPolicy.replace |
Перенести данные в новую Job в хвосте и отменить прежнюю группу. |
Debounce и throttle
Заголовок раздела «Debounce и throttle»final class Search extends Solo<SearchState> { final SearchApi api; late final _queries = accumulate<SearchState, String, void>( (ctx, text) async { final results = await ctx.wait(() => api.search(text)); ctx.emit(SearchState.results(results)); }, merge: (previous, incoming) => incoming, timing: AccumulationTiming.debounce(const Duration(milliseconds: 300)), policy: AccumulationPolicy.join, key: 'query', );
Search(this.api) : super(const SearchState.idle());
SoloJob<void> query(String text) => _queries.add(text);}Параметр timing откладывает готовность группы к запуску:
AccumulationTiming.debounce(duration)ждёт паузы во входных событиях.AccumulationTiming.throttle(duration)ограничивает частоту запуска групп, отсчитывая интервал от фактического старта предыдущей группы.
Эти задержки не занимают место работающей Job. Ожидающие группы
остаются в queue и пропускают другие готовые Job. Сами обработчики
по-прежнему выполняются по одному. Ожидание времени не отбрасывает
данные: collect хранит принятые события, а accumulate хранит результат
merge. Без timing или с Duration.zero группы готовы сразу.
Этот контроллер поиска сохраняет последний запрос и ждёт 300 мс без
нового ввода перед запуском. SearchApi и SearchState являются типами
приложения:
Уже начавшийся поиск завершается до старта следующего. Возвращённая Job предоставляет исход и отмену, как и другие Job. Закрытие отменяет ожидающие группы, а не досылает их. Разделы ниже разбирают группировку, порядок выполнения и исходы отдельных вызовов.
Объединение правок настроек
Заголовок раздела «Объединение правок настроек»Правка описывает, какие настройки изменить. Новое значение поля заменяет прежнее; отсутствующие в новой правке поля сохраняются. В этом примере сами настройки не могут быть null, поэтому null в правке означает «оставить поле без изменения». Для nullable-настройки с возможностью очистки понадобится отдельно представлять наличие значения в правке.
import 'package:solo/solo.dart';
class Settings { final bool notifications; final String theme; final String language;
const Settings({ required this.notifications, required this.theme, required this.language, });}
class SettingsPatch { final bool? notifications; final String? theme; final String? language;
const SettingsPatch({ this.notifications, this.theme, this.language, });
SettingsPatch merge(SettingsPatch incoming) => SettingsPatch( notifications: incoming.notifications ?? notifications, theme: incoming.theme ?? theme, language: incoming.language ?? language, );
Settings apply(Settings current) => Settings( notifications: notifications ?? current.notifications, theme: theme ?? current.theme, language: language ?? current.language, );}
abstract interface class SettingsApi { Future<void> save(Settings settings);}
class SettingsController extends Solo<Settings> { final SettingsApi _api; late final _updates = accumulate<Settings, SettingsPatch, void>( (ctx, patch) async { final next = patch.apply(ctx.state); await ctx.join(() => _api.save(next)); ctx.emit(next); }, merge: (accumulated, incoming) => accumulated.merge(incoming), key: 'settings', timing: AccumulationTiming.debounce(const Duration(milliseconds: 200)), );
SettingsController(this._api, Settings initial) : super(initial);
SoloJob<void> update(SettingsPatch patch) => _updates.add(patch);}Первое событие становится накопленным значением без вызова merge.
Каждое следующее синхронно вызывает merge(accumulated, incoming) из
add, до старта обработчика. Сохраняется только результат. false —
заданное значение, которое функция слияния в примере сохраняет.
Например, внесите три изменения в одну группу:
Future<void> changeSettings(SettingsApi api, Settings initial) async { final controller = SettingsController(api, initial); try { final first = controller.update( const SettingsPatch(notifications: true), ); final second = controller.update(const SettingsPatch(theme: 'dark')); final third = controller.update(const SettingsPatch(language: 'ru'));
final outcomes = await Future.wait([ first.done, second.done, third.done, ]); for (final outcome in outcomes) { print(outcome); } } finally { await controller.close(); }}При политике по умолчанию вызовы разделяют один SoloJob<void>.
Его обработчик сохраняет один снимок со всеми тремя правками, затем
излучает этот снимок после 200 мс без новых правок. До запуска обработчика
состояние остаётся исходным. За время паузы могут выполниться готовые Job.
Наблюдение .done сообщает Done, Failed или Cancelled без броска.
Если сохранение падает, состояние не меняется, а группа завершается
провалом. Повторять неудачные правки решает вызывающий; более поздняя
независимая правка не включает их автоматически. Пример предполагает
одного писателя и future API, которая завершается при фактическом
окончании записи. join удерживает слот очереди до её завершения,
в том числе после отмены. Один клиентский таймаут не доказывает, что
сервер перестал записывать. Отмена также может помешать финальному
emit после того, как сервер уже принял запись.
Команды, где важна только последняя
Заголовок раздела «Команды, где важна только последняя»Очередь, в которой стоят resume, а за ним pause, собирается сделать
две взаимно уничтожающие вещи, и пришедший сейчас resume обесценивает
всю пару: её итог — то, к чему привёл бы один resume. Отменять ничего
не нужно, потому что накопитель не создаёт задач, которые пришлось бы
отменять. merge берёт пришедшую команду и выбрасывает ту, что была:
import 'package:solo/solo.dart';
enum Command { resume, pause }
sealed class Playback { const Playback();}
final class Playing extends Playback { const Playing();}
final class Paused extends Playback { const Paused();}
abstract interface class PlayerDevice { Future<void> resume();
Future<void> pause();}
final class Player extends Solo<Playback> { final PlayerDevice device;
late final _transport = accumulate<Playback, Command, void>( key: 'transport', merge: (accumulated, incoming) => incoming, (ctx, command) async { switch (command) { case Command.resume: await ctx.join(device.resume); ctx.emit(const Playing()); case Command.pause: await ctx.join(device.pause); ctx.emit(const Paused()); } }, );
Player(this.device) : super(const Paused());
SoloJob<void> resume() => _transport.add(Command.resume);
SoloJob<void> pause() => _transport.add(Command.pause);}Три вызова подряд — resume, pause, resume — оставляют в очереди
одну задачу, и все три возвращают один и тот же handle. Устройству
сказано возобновить один раз. По дороге ничего не ставилось в очередь и
не отменялось.
Это работает потому, что команды абсолютные: каждая говорит, каким будет итоговое состояние, поэтому последняя и есть ответ. Команды, которые надстраиваются друг над другом — «ещё на десять секунд вперёд», — объединяются сложением, а не заменой.
adjacent, политика по умолчанию, присоединяет группу только к хвосту
очереди, поэтому задача другого рода, добавленная между двумя командами,
становится границей и склейка на ней прекращается. Именно это сохраняет
порядок с остальной работой; AccumulationPolicy.join от него
отказывается и присоединяет группу там, где она стоит.
Когда это всё-таки разные задачи
Заголовок раздела «Когда это всё-таки разные задачи»Там, где resume и pause правда разные операции, а не одна команда со
значением, они обычные задачи со своими ключами, и Policy.replace не
поможет: она убирает задачи с тем же ключом, а общего ключа у этих двух
нет. Убирает их собственная очередь контроллера:
SoloJob<void> resume() { queue.removeWhere( (job) => job.key == Command.resume || job.key == Command.pause, );
return run<Playback, void>(key: Command.resume, (ctx) async { await ctx.join(device.resume); ctx.emit(const Playing()); });}Работающую задачу очередь не трогает никогда: уже начавшийся pause
доработает до конца, что бы за ним ни убрали. Для него нужен
cancelAll() — он чистит очередь и отменяет текущую — или общий ключ на
обе команды с Policy.restart. Задачи, созданные с cancellable: false,
пропускаются, пока не передашь force: true.
Сбор записей журнала
Заголовок раздела «Сбор записей журнала»Для логов важны отдельные записи и их порядок. collect дописывает
события в приватный список и передаёт обработчику неизменяемый снимок
при закрытии состава группы. Повторяющиеся события сохраняются, весь список при
каждом добавлении не копируется. Сами записи не клонируются;
используйте неизменяемые объекты событий.
import 'package:solo/solo.dart';
class LogEntry { final String message;
const LogEntry(this.message);}
abstract interface class LogApi { Future<void> send(List<LogEntry> entries);}
class LogController extends Solo<int> { final LogApi _api; late final _logs = collect<int, LogEntry, void>( (ctx, entries) async { await ctx.join(() => _api.send(entries)); ctx.emit(ctx.state + entries.length); }, key: 'logs', policy: AccumulationPolicy.join, timing: AccumulationTiming.throttle(const Duration(seconds: 1)), );
LogController(this._api) : super(0);
SoloJob<void> logEvent(LogEntry entry) => _logs.add(entry);}Три вызова до выполнения отправят один список из трёх записей. Первая
группа готова сразу; следующие стартуют с интервалом не менее секунды.
Пришедшие за этот интервал записи сохраняются для следующей группы. Состояние
этого контроллера считает записи, отправка которых завершилась,
а обработчик дошёл до emit. Раздел о политиках объясняет выбор join
в качестве правила аккумуляции в этом примере.
Сбор записей не гарантирует доставку. Неудачная отправка, отмена группы или закрытие контроллера могут оставить их неотправленными. Хранение на диске, повторы и финальная отправка при закрытии требуют протокола приложения сверх этого накопителя.
Выбор момента готовности группы
Заголовок раздела «Выбор момента готовности группы»// Пауза во вводе: каждый add перезапускает таймер на 200 мс.late final _queries = accumulate<Ready, String, void>( (ctx, text) => ctx.wait(() => api.search(text)), merge: (previous, incoming) => incoming, timing: AccumulationTiming.debounce(const Duration(milliseconds: 200)),);
// Потолок частоты запусков, отсчитываемый от фактического старта.late final _metrics = collect<Ready, Metric, void>( (ctx, events) => ctx.join(() => api.send(events)), timing: AccumulationTiming.throttle(const Duration(seconds: 5)),);Обе фабрики принимают необязательный timing. Настройка определяет,
когда группе можно стартовать; collect сохраняет каждое принятое событие,
а accumulate — результат своей merge. Например,
merge: (previous, incoming) => incoming сохраняет последнее значение.
AccumulationTiming.debounce(duration) ждёт паузы после последнего
принятого события в каждой группе. Каждое добавление перезапускает таймер,
даже если merge вернула неизменное значение. При интервале 200 мс
события на 0, 60 и 120 мс делают группу готовой на 320 мс. Непрерывный
вход может удерживать открытую группу неограниченно долго. Когда таймер
срабатывает, состав группы закрывается, даже если другой Job работает.
Следующие события создают новую группу и не меняют закрытый вход.
collect создаёт неизменяемый снимок один раз, при закрытии состава.
AccumulationTiming.throttle(duration) разрешает первый старт сразу,
затем выдерживает не менее duration между фактическими стартами групп
одного накопителя. Добавления за этот интервал собираются в открытую
группу, не продлевая таймер. По истечении интервала накопленное готово
без нового события. Группа принимает события до извлечения из очереди
для выполнения. Для пустого интервала Job не создаётся.
Старт — переход в running, до onStart. Группа, отклонённая стартовыми
правилами, не расходует интервал throttle; отмена из onStart расходует.
Если группа стартовала на 0 мс с интервалом 200 мс, а другой Job занимает
слот до 500 мс, следующая группа стартует на 500 мс, а последующая —
не раньше 700 мс. Если собственный обработчик группы, дети или уборка
длятся дольше интервала, следующая группа может стартовать сразу после
их завершения.
Ожидающие группы видны в queue.jobs. Очередь берёт первый готовый Job
в порядке списка, пропуская ожидающие группы. Поэтому список может быть
непустым, когда ни один Job не работает. Ожидание времени не занимает
слот выполнения; после старта обработчика очередь ждёт его детей
и уборку перед запуском следующего корневого Job.
Без timing или при Duration.zero группы готовы сразу и таймеров нет.
Отрицательная длительность бросает ArgumentError. Общую настройку
времени можно передать нескольким накопителям; у каждого свои группы
и интервал throttle. Колбэки таймеров следуют порядку событий Dart:
событие, обработанное до колбэка debounce, может перезапустить таймер;
событие после него относится к новой группе. Занятый цикл событий
может задержать колбэки и старты.
Выбор места присоединения событий
Заголовок раздела «Выбор места присоединения событий»// Все три, пока текущая Job ещё удерживает очередь:settings.change(a1); // A1, в этот накопительsettings.save(); // B, обычная отдельная Jobsettings.change(a2); // A2 — куда он ляжет, решает политикаОбе фабрики принимают AccumulationPolicy, фиксируемую при создании
накопителя. Рассмотрим события A1, B, A2, добавленные, пока текущая
задача удерживает очередь. A принадлежит одному накопителю; B — другая
задача.
| Правило | Очередь после добавлений | Хэндл для A2 |
|---|---|---|
adjacent (по умолчанию) |
[A1, B, A2] |
Новая задача |
replace |
[B, A(A1 + A2)] |
Новая задача; задача A1 отменяется |
join |
[A(A1 + A2), B] |
Прежняя задача A1 |
adjacent принимает событие в группу, только если она стоит в хвосте
очереди. replace и join находят последнюю открытую ожидающую группу
того же накопителя, проходя через другие задачи. Эти задачи остаются в очереди.
При join A сохраняет место перед B; replace переносит A за B.
Выполнение также зависит от готовности: готовая B может пройти перед A,
пока A ждёт свой временной срок.
Правила используют текущую очередь. Если B уже выполнилась, ожидающая A
снова может оказаться хвостом, и следующее событие присоединится к ней
при adjacent. Уже раздельные группы задним числом не объединяются.
Закрытая debounce-группа не подходит ни одному из трёх правил.
replace переносит все накопленные данные в новую ожидающую задачу
и добавляет входящее событие. Новый хэндл создаётся даже тогда, когда
старая группа уже была хвостом. Старая задача завершается Cancelled
с ManualCancelReason, started: false и описанием
replaced by accumulated group. Новая группа находится в очереди до
вызова обработчиков завершения и отмены прежней задачи. Эти обработчики
могут снова добавить, отменить или закрыть, поэтому возвращённый новый
хэндл может быть уже отменён к моменту возврата внешнего add. Замена
начинает новый интервал debounce и сохраняет интервал throttle накопителя.
Замена ожидающей группы обязательна, в том числе при настройке
cancellable: false. Работающую группу она никогда не отменяет.
У обычной политики очереди Policy.replace свои прежние правила отмены;
AccumulationPolicy — отдельный enum.
Возможность объединения определяется идентичностью накопителя. Два
накопителя с одинаковым key остаются разными. Ключ по-прежнему помечает
их задачи для обычного поиска и политик очереди. Создание накопителя
заново на каждом вызове не даёт событиям присоединяться к существующему.
Запуск, отмена и ошибки
Заголовок раздела «Запуск, отмена и ошибки»final group = settings.report(metric);
// Все, кто добавлял в эту группу, держат одну и ту же ручку...switch (await group.done) { case Done(): print('sent'); case Failed(:final error): print('failed: $error'); case Cancelled(:final reason): print('cancelled: $reason');}
// ...и её отмена отменяет всю группу.await group.cancel();Все три правила работают с ожидающими группами. Группа прекращает приём
событий по истечении debounce или при извлечении из очереди,
до canStart и onStart. События
из этих колбэков идут в последующую группу. Обработчик получает заданный
рабочий тип и правила, а очередь ждёт его детей и уборку так же,
как у любого другого SoloJob.
До первого события накопитель не создаёт задач. Без timing, если каждая
группа завершается до прихода следующего события,
каждое событие запускает отдельную задачу. Поиск при join может пройти
очередь; adjacent проверяет только её хвост. replace также ищет
по всей очереди.
При adjacent и join добавления в одну группу разделяют один хэндл,
результат и отмену. Отмена хэндла затрагивает всю группу. При replace
старые хэндлы остаются отменёнными: они не перенаправляются на новую
задачу, и отмена уже заменённого хэндла не отменяет его преемника.
Ожидайте новый хэндл, чтобы наблюдать перенесённую работу.
merge должна быть синхронной и чистой. Она не должна менять аргументы
или обращаться к контроллеру. Если она бросает, add бросает ту же
ошибку, а существующая группа и таймер остаются неизменными; даже replace
сохраняет старую задачу. Вызов add того же накопителя из его merge
бросает StateError. Движок также проверяет после колбэка, что группа
ещё подходит, прежде чем записать результат.
canStart, keepWhile и отмена действуют на всю задачу. Стартовое
правило может отменить весь её накопленный вход. Результат и ошибки
обработчика следуют обычному контракту Job; принятая отмена сохраняет
приоритет над более поздним значением или ошибкой. Накопитель не
откатывает частичные внешние эффекты и не пересылает неудачный вход сам.
Отмена ожидающей debounce-группы убирает её таймер. Удаление или очистка throttle-групп сохраняет уже начатый интервал, поэтому новое событие не обходит ограничение. Таймер истекает один раз и не возобновляется до старта очередной группы.
close() отменяет все таймеры timing, снимает ожидающие группы и отменяет
или ждёт работающую задачу по обычным правилам. Финальный пакет он не
отправляет. Вызов add после
начала закрытия возвращает новую задачу, уже завершённую
Cancelled(closed); merge и обработчик не вызываются. После завершения
группы её внутреннее хранилище входа освобождается. Обработчик, результат
или ошибка, сохранившие вход, по-прежнему владеют своими ссылками.