Typed DAG task workflow library with Pydantic models, lifecycle hooks, and fail-fast semantics.
source .venv/bin/activate # activate venv
pip install -e ".[dev]" # install with dev deps
pytest -v # run tests
ruff check taskmaestro/ tests/ examples/ # lint
ruff format taskmaestro/ tests/ examples/ # format
mypy taskmaestro # type check (strict mode)| File | Responsibility |
|---|---|
taskmaestro/exceptions.py |
Exception hierarchy (no internal deps) |
taskmaestro/context.py |
ExecutionContext with correlation ID, logger, scratch dir, service registry |
taskmaestro/task.py |
Task[I, O] ABC, type introspection (get_input_type, get_output_type) |
taskmaestro/workflow.py |
Workflow (linear + DAG), WorkflowBuilder, validation (cycles, types, fan-in) |
taskmaestro/job.py |
Job[C], JobStatus, TaskStatus, TaskResult dataclass |
taskmaestro/runner.py |
Runner — topological execution, timeout via signal.alarm, hook dispatch |
taskmaestro/hooks/base.py |
Event StrEnum, Hook protocol, BaseHook no-op base |
taskmaestro/hooks/logging.py |
LoggingHook — logs events via logging module |
taskmaestro/hooks/timing.py |
TimingHook — records durations via time.monotonic() |
taskmaestro/hooks/persistence.py |
ResultPersistenceHook — writes {task_name}.json per task |
- Type introspection: Walk MRO via
__orig_bases__+typing.get_args()to extract concreteI/Otypes - Fan-in: Downstream task input model fields mapped to upstream outputs via
model_fields(Pydantic v2) - Timeouts:
signal.alarm(Unix only, main thread); gracefully warns if unavailable - Hook error swallowing:
_emit()wraps each hook call in try/except, reports viawarnings.warn() - Validation order: unique names → acyclic (DFS) → type chain → result task detection
- Shared fixtures and reusable tasks/models in
tests/conftest.py - Tests organized by module:
test_exceptions,test_context,test_task,test_workflow,test_job,test_runner,test_hooks - Timeout tests skip on non-Unix (no
signal.SIGALRM) - Use
RecordingHookpattern to assert event sequences