A Python 3.12+ library for defining and executing typed DAG task workflows with Pydantic models, lifecycle hooks, and fail-fast semantics.
python3 -m venv .venv
source .venv/bin/activate
pip install -e ".[dev]"from pydantic import BaseModel
from taskmaestro import Task, Workflow, Job, Runner, ExecutionContext
class NumberInput(BaseModel):
value: int
class NumberOutput(BaseModel):
value: int
class AddOne(Task[NumberInput, NumberOutput]):
def run(self, input: NumberInput, ctx: ExecutionContext) -> NumberOutput:
return NumberOutput(value=input.value + 1)
class Double(Task[NumberOutput, NumberOutput]):
def run(self, input: NumberOutput, ctx: ExecutionContext) -> NumberOutput:
return NumberOutput(value=input.value * 2)
workflow = Workflow(name="math", tasks=[AddOne, Double])
job = Job(workflow=workflow, config=NumberInput(value=5))
result = Runner().run(job)
print(result.status) # "completed"
print(result.result.value) # 12from pydantic import BaseModel
from taskmaestro import Task, Workflow, Job, Runner, ExecutionContext
class Input(BaseModel):
value: int
class Output(BaseModel):
value: int
class MergedInput(BaseModel):
a: Output
b: Output
class MergedOutput(BaseModel):
total: int
class BranchA(Task[Input, Output]):
def run(self, input: Input, ctx: ExecutionContext) -> Output:
return Output(value=input.value + 1)
class BranchB(Task[Input, Output]):
def run(self, input: Input, ctx: ExecutionContext) -> Output:
return Output(value=input.value * 2)
class Merge(Task[MergedInput, MergedOutput]):
def run(self, input: MergedInput, ctx: ExecutionContext) -> MergedOutput:
return MergedOutput(total=input.a.value + input.b.value)
workflow = (
Workflow.builder(name="fan_in")
.add_task(BranchA)
.add_task(BranchB)
.add_task(Merge, depends_on={"a": BranchA, "b": BranchB})
.build()
)
job = Job(workflow=workflow, config=Input(value=5))
result = Runner().run(job)
print(result.status) # "completed"
print(result.result.total) # 16 (6 + 10)You define Tasks (typed units of work), compose them into a Workflow (linear chain or DAG), bind input data via a Job, and hand it to a Runner for execution. Type safety is enforced at build time — input/output models are validated across the entire graph. Fail-fast semantics stop execution on the first error, and Hooks provide cross-cutting lifecycle observations without coupling to task logic.
| Concept | Description |
|---|---|
| Task | Subclass Task[I, O] with Pydantic models for input and output, then implement run(input, ctx). Each task can declare an optional timeout_seconds. For tasks with multiple named outputs, use inline Inputs/Outputs classes inside the task body. |
| Workflow | Build a linear pipeline with Workflow(tasks=[...]) or a DAG with Workflow.builder(). Prefer builder.task() and task handles for unambiguous dependencies; the fluent add_task() API remains supported. The builder accepts collect() for gathering outputs into collection fields and mapped_over=TaskMap(...) for sequential expansion over configured mappings. Use config_fields to declare which input fields come from JobConfiguration. Workflows are validated at build time for cycles, type compatibility, and input completeness. |
| Job | Binds a Workflow to a typed config (the root task's input). Tracks status (pending → running → completed/failed), the final result, any error, and per-task task_results. Optionally accepts a JobConfiguration for per-task static config values. A job can only be run once. |
| Runner | Executes tasks in topological order, stopping on the first failure (fail-fast). Supports per-task and per-job timeouts via signal.alarm (Unix only). Dispatches lifecycle events to registered hooks. |
| ExecutionContext | Passed to every run() call. Provides a logger, an auto-generated correlation_id (UUID), a scratch_dir (temporary directory), and a service registry (register()/resolve()) for injecting shared resources like DB connections. |
| Hooks | Subclass BaseHook and override methods like on_job_start, on_task_complete, etc. Hook errors are swallowed and reported via warnings.warn(), so they never crash the job. Built-ins: LoggingHook, TimingHook, ResultPersistenceHook. |
| ObjectModel | Generic ObjectModel[T] base model for wrapping arbitrary (non-Pydantic) objects. Enables arbitrary_types_allowed so fields can hold native library objects like database connections or API clients. |
For common cases, run a workflow directly without constructing Job and Runner:
result = workflow.run(
Input(value=5),
task_config={"configured_task": {"option": "value"}},
hooks=[LoggingHook()],
)The explicit Job and Runner API remains available for advanced lifecycle control.
builder.task() adds a task and returns a handle to that specific instance. Handles avoid ambiguous class and string references, especially when the same task class is registered more than once:
builder = Workflow.builder(name="parallel_wells")
model = builder.task(LoadModel)
well_1 = builder.task(LoadWellPath, name="well_1", depends_on=model)
well_2 = builder.task(LoadWellPath, name="well_2", depends_on=model)
builder.task(Process, name="proc_1", depends_on=well_1)
builder.task(Process, name="proc_2", depends_on=well_2)
workflow = builder.build()Use handle.field("field_name") to route one output field. Named input dependencies can be passed directly to task(), while keyword arguments to collect() provide a concise keyed collection:
merged = builder.task(
MergeResults,
primary=producer.field("result"),
checks=collect(tests=tests, lint=lint, types=types),
)
builder.set_result_task(merged)Handles are accepted anywhere dependency references are accepted. A handle from a different builder is rejected. Use the depends_on dictionary form when an input field conflicts with a reserved builder argument such as name or config_fields. The existing fluent add_task() API remains fully supported for backward compatibility.
The same Task class can appear multiple times in a workflow with different names. With the fluent API, use the name= parameter in add_task():
workflow = (
Workflow.builder(name="parallel_wells")
.add_task(LoadModel)
.add_task(LoadWellPath, name="well_1", depends_on=LoadModel)
.add_task(LoadWellPath, name="well_2", depends_on=LoadModel)
.add_task(Process, name="proc_1", depends_on="well_1")
.add_task(Process, name="proc_2", depends_on="well_2")
.build()
)Named instances can be referenced as string dependencies (depends_on="well_1") or in fan-in dicts.
JobConfiguration provides static config values for individual tasks, merged with upstream outputs at runtime. Declare which fields come from config with config_fields:
from taskmaestro import EmptyConfig, Job, JobConfiguration, Workflow
workflow = (
Workflow.builder(name="configured")
.add_task(LoadModel, config_fields=["path"]) # root: all input from config
.add_task(Transform, depends_on=LoadModel, config_fields=["scale_factor"]) # mixed
.build()
)
job_config = JobConfiguration(
{
"load_model": {"path": "/data/model.egrid"},
"transform": {"scale_factor": 2.5},
}
)
job = Job(workflow=workflow, config=EmptyConfig(), job_configuration=job_config)
result = Runner().run(job)A workflow can be wrapped as a typed task and used inside a larger workflow. Its input type is inferred from its root task and its output type from its result task, so normal workflow type validation still applies at both boundaries:
inner = Workflow(name="normalize", tasks=[CleanText, NormalizeText])
Normalize = inner.as_task(name="normalize_text")
outer = (
Workflow.builder("document_pipeline")
.add_task(LoadDocument)
.add_task(Normalize, depends_on=LoadDocument)
.add_task(IndexDocument, depends_on=Normalize)
.build()
)workflow_task(inner, name="normalize_text") is the equivalent factory-style API.
The inner workflow must have exactly one root that receives input from the outer
workflow. Roots supplied entirely by JobConfiguration are excluded; if every root is
configured, pass that configuration to as_task() and the wrapper accepts
EmptyConfig:
ConfiguredPipeline = inner.as_task(
name="configured_pipeline",
job_configuration=inner_config,
)The inner tasks share the outer ExecutionContext, including services, scratch directory,
and correlation ID. The wrapper is an opaque lifecycle boundary: outer runner hooks and
Job.task_results see one wrapper task, while an inner failure is reported with the inner
workflow and failed task names. Mermaid visualization expands wrappers as subgraphs.
Route a specific field from an upstream task's output (rather than the whole output) using (Task, "field") tuples:
class ExtractKeywords(Task):
class Inputs(BaseModel):
content: TextContent
class Outputs(BaseModel):
keywords: KeywordsOutput
num_words_removed: int
def run(self, input: Inputs, ctx: ExecutionContext) -> Outputs: ...
workflow = (
Workflow.builder(name="analysis")
.add_task(ExtractKeywords, depends_on=PrepareText)
.add_task(
BuildReport,
depends_on={
"keywords": (ExtractKeywords, "keywords"), # routes .keywords field
"num_words_removed": (
ExtractKeywords,
"num_words_removed",
), # routes .num_words_removed
"stats": ComputeWordStats, # whole output
},
)
.build()
)Use collect() when several task outputs should populate one list[T] or
dict[str, T] field. Positional members preserve declaration order:
from taskmaestro import collect
class GridInput(BaseModel):
surfaces: list[Surface]
workflow = (
Workflow.builder("create_grid")
.add_task(LoadSurface, name="top")
.add_task(GenerateSurface, name="middle")
.add_task(LoadSurface, name="base")
.add_task(
CreateGrid,
depends_on={"surfaces": collect("top", "middle", "base")},
)
.build()
)Use a mapping to preserve aliases in a dict[str, T], and use (task, "field")
to collect a specific output field:
depends_on={
"surfaces": collect({
"top": ("top_loader", "surface"),
"base": ("base_loader", "surface"),
})
}The equivalent YAML forms are:
depends_on:
surfaces:
collect:
- top
- [middle, generated_surface]
- basedepends_on:
surfaces:
collect:
top: [top_loader, surface]
base: [base_loader, surface]Every member is checked against the field's element type when the workflow is
built. Subtypes are accepted. collect() and collect({}) explicitly create
empty list and dictionary inputs, respectively.
A mapped task invokes one task declaration for every entry in a configured
mapping. Mapped items execute sequentially in mapping declaration order.
Each item gets a fresh task instance and child ExecutionContext.
builder = Workflow.builder("create_grid")
connection = builder.task(ConnectToResInsight)
surfaces = builder.map_task(
LoadRegularSurface,
name="load_surfaces",
depends_on={"resinsight": connection},
config_fields=["unit"],
over="surfaces",
key_as="surface_name",
value_as="path",
error_mode="fail_fast",
)
builder.task(CreateGrid, depends_on={"surfaces": surfaces})
workflow = builder.build()The mapped task's input model contains the injected key and value fields, not the source mapping:
class LoadSurfaceInput(BaseModel):
resinsight: RipsInstance
unit: str
surface_name: str # key_as
path: str # value_asConfigure the source through JobConfiguration:
job_configuration = JobConfiguration({
"load_surfaces": {
"unit": "meters",
"surfaces": {
"top": "/data/top.irap",
"base": "/data/base.irap",
},
},
})The logical output is a MappedOutput[O] Pydantic root model containing an
insertion-ordered dict[str, O], where O is the task's declared output type.
When a mapped task is connected to a named dict[str, O] input field, its
root value is unwrapped automatically:
class CreateGridInput(BaseModel):
surfaces: dict[str, RegularSurface]The equivalent YAML task declaration is:
- task: resinsight.load_regular_surface
name: load_surfaces
map:
over: surfaces
key_as: surface_name
value_as: path
error_mode: fail_fast
depends_on:
resinsight: resinsight.connect
config_fields: [unit]Input YAML:
load_surfaces:
unit: meters
surfaces:
top: /data/top.irap
base: /data/base.irapfail_fast stops at the first failed item. collect_all attempts every item
and reports an aggregate MappedTaskExecutionError. An empty mapping succeeds
with MappedOutput(root={}). Per-item records are available in
job.mapped_item_results, and
built-in logging, timing, and persistence hooks observe individual items.
Concurrent mapped execution is intentionally deferred.
ObjectModel[T] wraps arbitrary (non-Pydantic) objects so they can flow through workflows. Use it as a type alias for simple wrappers, or subclass it to add extra fields:
from taskmaestro import ObjectModel
# Type alias — no extra fields needed
GridCase = ObjectModel[rips.EclipseCase]
WellPath = ObjectModel[rips.WellPath]
# Subclass — adds fields alongside the wrapped object
class AddPerforationInput(ObjectModel[rips.WellPath]):
start_md: float
end_md: float
# Access the wrapped object via .value
grid = GridCase(value=eclipse_case)
print(grid.value.name)Installed packages can publish tasks and workflows using standard Python entry points:
[project.entry-points."taskmaestro.tasks"]
"acme.prepare" = "acme_tasks.prepare:Prepare"
[project.entry-points."taskmaestro.workflows"]
"acme.analysis" = "acme_tasks.workflows:analysis_workflow"A task entry point must resolve to a Task subclass and a workflow entry point must
resolve to a Workflow instance. Prefix names with the provider name to avoid clashes.
Consumers can discover plugins without scanning package directories:
from taskmaestro import registered_tasks, registered_workflows
tasks = registered_tasks() # dict[str, type[Task]]
workflows = registered_workflows() # dict[str, Workflow]Use registered_task_names() and registered_workflow_names() to inspect identifiers
without importing plugin modules, or get_registered_task(name) and
get_registered_workflow(name) to load one plugin. Duplicate names and invalid plugin
types raise PluginLoadError.
Workflows can be defined entirely in YAML instead of Python. A task: value may be
either a registered task identifier or a dotted Python class path. The loader resolves
registered identifiers first and validates the full configuration:
# workflow.yaml
workflow:
name: text_analysis
tasks:
- task: pipeline.PrepareText
- task: pipeline.GenerateStopWords
- task: pipeline.ComputeWordStats
depends_on:
content: pipeline.PrepareText
stop_words: pipeline.GenerateStopWords
- task: pipeline.ExtractKeywords
depends_on:
content: pipeline.PrepareText
stop_words: pipeline.GenerateStopWords
- task: pipeline.ScoreReadability
depends_on: pipeline.PrepareText
- task: pipeline.BuildReport
depends_on:
stats: pipeline.ComputeWordStats
keywords: [pipeline.ExtractKeywords, keywords] # output field routing
readability: pipeline.ScoreReadability
num_words_removed: [pipeline.ExtractKeywords, num_words_removed] # output field routing
runner:
hooks:
- hook: taskmaestro.hooks.logging.LoggingHook
- hook: taskmaestro.hooks.timing.TimingHook
context:
services:
title: "Python Overview"# input.yaml
prepare_text:
text: "Python is a high-level programming language..."
title: "Python Overview"Load and run:
from taskmaestro import load_workflow_from_yaml, run_workflow_from_yaml
# Load for inspection, then run
loaded = load_workflow_from_yaml("workflow.yaml", "input.yaml")
result = loaded.run()
# Or run directly
result = run_workflow_from_yaml("workflow.yaml", "input.yaml")YAML input always uses per-task configuration: every top-level key in input.yaml must be a registered task instance name, and its value must be a mapping or null. Unknown task names and scalar task values are rejected. Fields are validated against the task's input model and can configure root tasks, downstream tasks, and mapped tasks.
YAML also supports named task instances (name:), fan-in dictionaries, and output field routing via [task, field] lists. Named instances use their instance name as the input key:
load_well_path_1:
path: first.dev
load_well_path_2:
path: second.devWhen the same task class (or the same inner YAML file) appears more than once under different name:s, depends_on and result_task must use the instance name — referencing the class path is rejected as ambiguous.
Use workflow: instead of task: to compose another YAML workflow. Paths are resolved
relative to the containing workflow file, and workflow_input: optionally supplies the
inner workflow's per-task configuration:
workflow:
name: document_pipeline
tasks:
- task: pipeline.LoadDocument
- workflow: normalize/workflow.yaml
workflow_input: normalize/input.yaml
name: normalize_text
depends_on: pipeline.LoadDocument
- task: pipeline.IndexDocument
depends_on: normalize_textThe same root/result type inference and single-unconfigured-root requirement apply as for
Workflow.as_task().
Generate Mermaid diagrams of workflow topology:
print(workflow.to_mermaid())
# or with config nodes:
print(workflow.to_mermaid(job_configuration=job_config))Output:
---
title: text_analysis
---
graph TD
_start_(("start"))
_end_(("end"))
prepare_text["prepare_text"]
generate_stop_words["generate_stop_words"]
compute_word_stats["compute_word_stats"]
build_report["build_report"]
_start_ -->|TextInput| prepare_text
_start_ -->|TextInput| generate_stop_words
prepare_text -->|content: TextContent| compute_word_stats
compute_word_stats -->|WordStatsOutput| build_report
build_report -->|AnalysisReport| _end_
Edges are labeled with data types. Fan-in edges show field names, and field routing edges show .field: Type. When a JobConfiguration is provided, configured tasks get dashed edges from a JobConfiguration node.
WorkflowRunnerError (base)
├── WorkflowDefinitionError # Invalid workflow definition
│ ├── CycleDetectedError # Dependency cycle
│ └── IncompleteInputError # Missing fan-in field mappings
├── JobStateError # e.g., re-running a completed job
├── ConfigLoadError # YAML config loading failure
└── TaskExecutionError # Runtime task failure
├── MappedTaskExecutionError # One or more mapped items failed
├── TaskOutputTypeError # Output type mismatch
└── TaskTimeoutError # Task exceeded timeout
Installed packages provide a taskmaestro command for YAML workflows:
taskmaestro validate workflow.yaml --input input.yaml
taskmaestro graph workflow.yaml --input input.yaml
taskmaestro run workflow.yaml --input input.yaml --log-level INFOrun prints the final output as JSON and returns a nonzero exit code when the workflow fails. graph prints Mermaid markup.
Four full example pipelines are included in the examples/ directory:
| Example | Features |
|---|---|
examples/text_analysis/ |
DAG with fan-out/fan-in, output field routing, inline Inputs/Outputs classes, YAML config, Mermaid visualization |
examples/resinsight/ |
ObjectModel[T] for gRPC objects, JobConfiguration with per-task config, named task instances, config_fields, YAML config |
examples/image_processing/ |
Nested workflows through Workflow.as_task() and YAML workflow:, typed boundaries, expanded Mermaid subgraph |
examples/release_pipeline/ |
Keyed collect() dependencies, mapped tasks, mapped output routing, per-task YAML config |
Run an example:
python examples/text_analysis/pipeline.py # Python API
python examples/text_analysis/pipeline.py --yaml # YAML configsource .venv/bin/activate
pytest -v # run tests
ruff check . # lint
ruff format . # format
mypy taskmaestro # type check (strict)