Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
039d25c
Add collection dependencies for task outputs
kriben Sep 21, 2026
642b060
Add sequential mapped task expansion
kriben Sep 21, 2026
5f84df5
Validate generic container annotations recursively
kriben Sep 21, 2026
59eda9e
Fix job deadline, non-main-thread, and sub-second timeouts in Runner
kriben Sep 21, 2026
5436f40
Fix persistence path traversal, linear-mode YAML wiring, recursive refs
kriben Sep 21, 2026
40b94a6
Tighten Workflow validation: result_task, duplicates, empty list, edg…
kriben Sep 21, 2026
c495013
Preserve error detail from hooks and inner workflows
kriben Sep 21, 2026
bd90ee6
Reject ambiguous YAML task references instead of picking the last one
kriben Sep 21, 2026
695713a
Add workflow.input_mode and refuse to guess on ambiguous input YAML
kriben Sep 21, 2026
0c03911
Harden GitHub Actions workflows and align packaging metadata
kriben Sep 21, 2026
1623aaf
Add mapped release pipeline example
kriben Sep 28, 2026
7e7b963
Add backward-compatible task handles
kriben Sep 28, 2026
96a608b
Add Workflow.run convenience API
kriben Sep 28, 2026
b6f2c8b
Automatically unwrap mapped task outputs
kriben Sep 28, 2026
acd7d53
Add map_task builder convenience
kriben Sep 28, 2026
a28a556
Add keyword dependency binding
kriben Sep 28, 2026
7de1fd9
Add Taskmaestro command-line interface
kriben Sep 28, 2026
edb642b
Resolve CLI task imports relative to workflow
kriben Sep 28, 2026
2efddb8
Hide orphan start node in workflow graphs
kriben Sep 28, 2026
01fddd3
Require per-task YAML input
kriben Sep 28, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 15 additions & 3 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,19 @@ on:
branches: [main]
pull_request:

# Least privilege: the CI job only reads the repository.
permissions:
contents: read

# Cancel superseded runs for the same ref (e.g. force-pushes to a PR).
concurrency:
group: ci-${{ github.workflow }}-${{ github.ref }}
cancel-in-progress: true

jobs:
ci:
runs-on: ubuntu-latest
timeout-minutes: 15
strategy:
fail-fast: false
matrix:
Expand All @@ -20,18 +30,20 @@ jobs:
uses: actions/setup-python@v5
with:
python-version: ${{ matrix.python-version }}
cache: pip
cache-dependency-path: pyproject.toml

- name: Install dependencies
run: pip install -e ".[dev]"

- name: Ruff check
run: ruff check taskmaestro/ tests/
run: ruff check taskmaestro/ tests/ examples/

- name: Ruff format check
run: ruff format --check taskmaestro/ tests/
run: ruff format --check taskmaestro/ tests/ examples/

- name: Mypy type check
run: mypy taskmaestro

- name: Run tests with coverage
run: pytest --cov=taskmaestro --cov-report=term-missing
run: pytest --cov=taskmaestro --cov-report=term-missing --cov-fail-under=100
12 changes: 11 additions & 1 deletion .github/workflows/publish.yml
Original file line number Diff line number Diff line change
Expand Up @@ -4,17 +4,24 @@ on:
push:
tags: ["v*"]

# Default to read-only; the publish job opts into id-token below.
permissions:
contents: read

jobs:
build:
name: Build distributions
runs-on: ubuntu-latest
timeout-minutes: 10
steps:
- uses: actions/checkout@v4

- name: Set up Python
uses: actions/setup-python@v5
with:
python-version: "3.12"
cache: pip
cache-dependency-path: pyproject.toml

- name: Install build
run: pip install build
Expand All @@ -32,6 +39,7 @@ jobs:
name: Publish to PyPI
needs: build
runs-on: ubuntu-latest
timeout-minutes: 10
environment: pypi
permissions:
id-token: write
Expand All @@ -43,4 +51,6 @@ jobs:
path: dist/

- name: Publish via trusted publishing
uses: pypa/gh-action-pypi-publish@release/v1
# Third-party action pinned to a full commit SHA (tag v1.14.2); the
# `release/v1` branch is mutable and could be repointed.
uses: pypa/gh-action-pypi-publish@dc37677b2e1c63e2034f94d8a5b11f265b73ba33 # v1.14.2
3 changes: 2 additions & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,8 @@ mypy taskmaestro # type check (strict mode)
- **Type introspection**: Walk MRO via `__orig_bases__` + `typing.get_args()` to extract concrete `I`/`O` types
- **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 via `warnings.warn()`
- **Hook error swallowing**: `_emit()` wraps each hook call in try/except, reports via `warnings.warn(..., HookError, source=exc)` — message includes `repr(exc)`; `HookError` subclasses `UserWarning` so it can be filtered or escalated
- **Inner-workflow failures**: `workflow_task` raises `WorkflowTaskError` (a `TaskExecutionError`) carrying `inner_job` and chaining the original exception via `__cause__`; `Job.exception` keeps the raw exception alongside `Job.error`
- **Validation order**: unique names → acyclic (DFS) → type chain → result task detection

## Testing Conventions
Expand Down
226 changes: 220 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -110,16 +110,55 @@ You define **Tasks** (typed units of work), compose them into a **Workflow** (li
| 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()`. The builder accepts `depends_on` for single dependencies, fan-in dicts (`{"field": UpstreamTask}`), and `(Task, "field")` tuples for output field routing. Use `config_fields` to declare which input fields come from `JobConfiguration`. Workflows are validated at build time for cycles, type compatibility, and input completeness. |
| **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`:

```python
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:

```python
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:

```python
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. Use the `name=` parameter in `add_task()`:
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()`:

```python
workflow = (
Expand Down Expand Up @@ -231,6 +270,155 @@ workflow = (
)
```

## Collecting Multiple Outputs

Use `collect()` when several task outputs should populate one `list[T]` or
`dict[str, T]` field. Positional members preserve declaration order:

```python
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:

```python
depends_on={
"surfaces": collect({
"top": ("top_loader", "surface"),
"base": ("base_loader", "surface"),
})
}
```

The equivalent YAML forms are:

```yaml
depends_on:
surfaces:
collect:
- top
- [middle, generated_surface]
- base
```

```yaml
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`.

```python
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:

```python
class LoadSurfaceInput(BaseModel):
resinsight: RipsInstance
unit: str
surface_name: str # key_as
path: str # value_as
```

Configure the source through `JobConfiguration`:

```python
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:

```python
class CreateGridInput(BaseModel):
surfaces: dict[str, RegularSurface]
```

The equivalent YAML task declaration is:

```yaml
- 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:

```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:
Expand Down Expand Up @@ -324,8 +512,9 @@ context:

```yaml
# input.yaml
text: "Python is a high-level programming language..."
title: "Python Overview"
prepare_text:
text: "Python is a high-level programming language..."
title: "Python Overview"
```

Load and run:
Expand All @@ -341,7 +530,18 @@ result = loaded.run()
result = run_workflow_from_yaml("workflow.yaml", "input.yaml")
```

YAML supports named task instances (`name:`), per-task input config (keyed by task name in the input file), fan-in dicts, and output field routing via `[task, field]` lists.
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:

```yaml
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
Expand Down Expand Up @@ -405,19 +605,33 @@ WorkflowRunnerError (base)
├── 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:

```bash
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

Three full example pipelines are included in the `examples/` directory:
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:

Expand Down
3 changes: 2 additions & 1 deletion examples/image_processing/input.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -3,4 +3,5 @@
# Run:
# python examples/image_processing/pipeline.py --yaml --input input.yaml

image_path: "taskmaestro.png"
load_image:
image_path: "taskmaestro.png"
Loading
Loading