Руководства
Создание обработчика задач на Dart
Запускаемый обработчик с типизированными командами, группировкой, политиками допуска, отклонением и наблюдаемым состоянием.
Создание обработчика задач на Dart
Создадим координатор с двумя командами. SubmitTask принимает не более двух задач одной очереди за секунду. ReserveSlot ожидает короткое окно Debounce и отклоняет конфликтующее резервирование той же очереди работников.
Что вы изучите
- типизированные команды;
- группировку через
executionKey; - Rate Limit и Debounce на этапе допуска;
- публикацию неизменяемого состояния;
- ожидание конечных снимков и корректное закрытие.
Создайте проект
dart create -t console usecase_forge_tasks
cd usecase_forge_tasks
dart pub add usecase_forge
PROJECT STRUCTURE
usecase_forge_tasks/
├── bin/
│ └── usecase_forge_tasks.dart
├── lib/
│ └── task_processing.dart
└── pubspec.yaml
Определите состояние
import 'package:usecase_forge/usecase_forge.dart';
final class TaskProcessingState {
const TaskProcessingState({
required this.processedCount,
required this.lastTask,
});
const TaskProcessingState.empty() : processedCount = 0, lastTask = '';
final int processedCount;
final String lastTask;
@override
bool operator ==(Object other) =>
other is TaskProcessingState &&
other.processedCount == processedCount &&
other.lastTask == lastTask;
@override
int get hashCode => Object.hash(processedCount, lastTask);
}
Равенство по значению позволяет выходной границе отличать изменение состояния от повторной публикации того же значения.
Опишите два намерения
final class SubmitTask extends UseCaseCommand {
const SubmitTask({required this.queue, required this.taskId});
final String queue;
final String taskId;
@override
Object get executionKey => queue;
}
final class ReserveSlot extends UseCaseCommand {
const ReserveSlot({required this.queue, required this.taskId});
final String queue;
final String taskId;
@override
Object get executionKey => queue;
}
Ключ группирует операции по очереди. Занятость email не должна блокировать workers.
Зарегистрируйте политики допуска
final class TaskProcessingUseCase extends UseCase<TaskProcessingState> {
TaskProcessingUseCase()
: super(initialState: const TaskProcessingState.empty()) {
registerCommand<SubmitTask>(
_process,
instructions: UseCaseInstructionOverrides(
input: UseCaseInputInstructionOverrides(
rateLimit: UseCaseRateLimitInstructions(
maxExecutions: 2,
duration: const Duration(seconds: 1),
),
),
),
);
registerCommand<ReserveSlot>(
_reserve,
instructions: UseCaseInstructionOverrides(
input: UseCaseInputInstructionOverrides(
debounce: UseCaseDebounceInstructions(
duration: const Duration(milliseconds: 30),
),
conflictPolicy: UseCaseInputConflictPolicy.rejectNew,
),
),
);
}
final rejections = <(String, UseCaseCommandRejectionReason)>[];
}
Обе политики решают судьбу входа и не усложняют код обработчика.
Публикуйте состояние
Future<void> _process(
SubmitTask command,
UseCaseExecutionContext<TaskProcessingState> context,
) async => _publishProcessed(context, command.taskId);
Future<void> _reserve(
ReserveSlot command,
UseCaseExecutionContext<TaskProcessingState> context,
) async => _publishProcessed(context, command.taskId);
void _publishProcessed(
UseCaseExecutionContext<TaskProcessingState> context,
String taskId,
) {
final latest = context.snapshot.state;
context.publish(TaskProcessingState(
processedCount: latest.processedCount + 1,
lastTask: taskId,
));
}
Зафиксируйте типизированное отклонение
@override
Future<void> onCommandRejected(
UseCaseExecutionEntry<UseCaseCommand> entry,
UseCaseCommandRejectionReason reason,
) async {
final taskId = switch (entry.command) {
SubmitTask command => command.taskId,
ReserveSlot command => command.taskId,
_ => 'unknown',
};
rejections.add((taskId, reason));
}
Отклонение допуска не является ошибкой обработчика: обработчик не запускался.
Запустите сценарий
Future<void> main() async {
final tasks = TaskProcessingUseCase();
final threeFinished = tasks.stream
.where((snapshot) => snapshot.phase == UseCaseExecutionPhase.finished)
.take(3)
.drain<void>();
tasks.add(const SubmitTask(queue: 'email', taskId: 'email-1'));
tasks.add(const SubmitTask(queue: 'email', taskId: 'email-2'));
tasks.add(const SubmitTask(queue: 'email', taskId: 'email-3'));
tasks.add(const ReserveSlot(queue: 'workers', taskId: 'slot-1'));
tasks.add(const ReserveSlot(queue: 'workers', taskId: 'slot-2'));
await threeFinished;
print('Processed: ${tasks.state.state.processedCount}');
await tasks.close();
}
Полный запускаемый вариант находится в example/task_processing.dart основного пакета и выполняется его тестами документационных примеров.