Skip to content

Repository files navigation

Taskmaestro

Taskmaestro

A Python 3.12+ library for defining and executing typed DAG task workflows with Pydantic models, lifecycle hooks, and fail-fast semantics.

Installation

python3 -m venv .venv
source .venv/bin/activate
pip install -e ".[dev]"

Quick Start

Linear Pipeline

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)  # 12

DAG Workflow (Fan-In)

from 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)

Core Concepts

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.

Task Handles

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.

Named Task Instances

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.

Per-Task Configuration

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)

Nested Workflows

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.

Output Field Routing

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()
)

Collecting Multiple Outputs

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]
      - base
depends_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.

Mapped Tasks

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_as

Configure 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.irap

fail_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

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)

Plugin discovery

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.

YAML Configuration

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.dev

When 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_text

The same root/result type inference and single-unconfigured-root requirement apply as for Workflow.as_task().

Visualization

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_
Loading

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.

Error Handling

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

Command-Line Interface

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 INFO

run prints the final output as JSON and returns a nonzero exit code when the workflow fails. graph prints Mermaid markup.

Examples

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 config

Development

source .venv/bin/activate
pytest -v                  # run tests
ruff check .               # lint
ruff format .              # format
mypy taskmaestro       # type check (strict)

About

Workflow manager for Python tasks for ResInsight

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages