Руководства

Создание обработчика задач на Dart

Запускаемый обработчик с типизированными командами, группировкой, политиками допуска, отклонением и наблюдаемым состоянием.

v0.1.0-dev.2Начальный уровеньusecase_forge

Создание обработчика задач на Dart

Создадим координатор с двумя командами. SubmitTask принимает не более двух задач одной очереди за секунду. ReserveSlot ожидает короткое окно Debounce и отклоняет конфликтующее резервирование той же очереди работников.

Что вы изучите

  • типизированные команды;
  • группировку через executionKey;
  • Rate Limit и Debounce на этапе допуска;
  • публикацию неизменяемого состояния;
  • ожидание конечных снимков и корректное закрытие.
01

Создайте проект

Терминал
shell
dart create -t console usecase_forge_tasks
cd usecase_forge_tasks
dart pub add usecase_forge

PROJECT STRUCTURE

text
usecase_forge_tasks/
├── bin/
│   └── usecase_forge_tasks.dart
├── lib/
│   └── task_processing.dart
└── pubspec.yaml
02

Определите состояние

lib/task_processing.dart
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);
}

Равенство по значению позволяет выходной границе отличать изменение состояния от повторной публикации того же значения.

03

Опишите два намерения

lib/task_processing.dart
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.

04

Зарегистрируйте политики допуска

lib/task_processing.dart
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)>[];
}

Обе политики решают судьбу входа и не усложняют код обработчика.

05

Публикуйте состояние

lib/task_processing.dart
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,
  ));
}
06

Зафиксируйте типизированное отклонение

lib/task_processing.dart
@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));
}

Отклонение допуска не является ошибкой обработчика: обработчик не запускался.

07

Запустите сценарий

bin/usecase_forge_tasks.dart
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 основного пакета и выполняется его тестами документационных примеров.

ARKTELOS

Инженерные системы для программного обеспечения, которое обязано выдерживать нагрузку.