Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
76e9ffd
Record input assembly failures as task failures
kriben Sep 28, 2026
30ac029
Keep outer timeouts armed across nested workflow runs
kriben Sep 28, 2026
8086c44
Require configuration for every task with config_fields
kriben Sep 28, 2026
be187f7
Reject config_fields on output-field dependencies
kriben Sep 28, 2026
bdddcaf
Disarm task timers as soon as the task returns
kriben Sep 28, 2026
8f9708b
Return an independent workflow from WorkflowBuilder.build()
kriben Sep 28, 2026
8d6ff6d
Wrap all task-module import failures in ConfigLoadError
kriben Sep 28, 2026
840b444
Report unhashable YAML keys as YAML errors
kriben Sep 28, 2026
42ecbee
Wrap all hook construction failures in ConfigLoadError
kriben Sep 28, 2026
4922d20
Report recursive workflow references in nesting order
kriben Sep 28, 2026
124ff80
Support running the CLI with python -m
kriben Sep 28, 2026
6342b2f
Enter nested-workflow subgraphs at the input root
kriben Sep 28, 2026
bf11406
Resolve the image_processing example's image path independently of cwd
kriben Sep 28, 2026
b046b14
Report failed runs in the image_processing and text_analysis examples
kriben Sep 28, 2026
8e97496
Add agent quickstart with runnable workflow example
kriben Sep 28, 2026
e21da82
Rename agent instructions to AGENTS.md
kriben Sep 28, 2026
3078e7a
Add CLI task discovery with JSON schemas
kriben Sep 28, 2026
4b4991c
Describe runtime-only objects in task schemas
kriben Sep 28, 2026
f2c9b5e
Add structured JSON results for CLI validate and run
kriben Sep 28, 2026
1f9771e
Add JSON workflow inspection CLI
kriben Sep 28, 2026
1c5c75e
Cover remaining CLI inspection paths
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
12 changes: 9 additions & 3 deletions CLAUDE.md → AGENTS.md
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
# CLAUDE.md
# AGENTS.md

Typed DAG task workflow library with Pydantic models, lifecycle hooks, and fail-fast semantics.

For a runnable CLI/YAML example and the validate → graph → run loop, see [docs/agent-quickstart.md](docs/agent-quickstart.md). Run its commands from the repository root.

## Commands

```bash
Expand All @@ -21,6 +23,10 @@ mypy taskmaestro # type check (strict mode)
| `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/dependencies.py`, `taskmaestro/mapping.py` | Task handles, output references, `collect()`, mapped-task configuration |
| `taskmaestro/workflow_task.py` | Nested workflows wrapped as tasks |
| `taskmaestro/yaml_config.py` | YAML parsing, task imports, workflow and input validation |
| `taskmaestro/cli.py`, `taskmaestro/discovery.py` | CLI (`run`, `validate`, `graph`, `tasks list/describe`, `workflow describe`), plugin entry-point discovery |
| `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 |
Expand All @@ -32,14 +38,14 @@ 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
- **Timeouts**: `signal.setitimer`/`SIGALRM` (Unix only, main thread); gracefully warns if unavailable. All armed timeouts go through the process-wide `_AlarmScheduler`, which always arms the nearest expiry, so nested runs (`workflow_task`) cannot cancel or extend an enclosing run's timeouts
- **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

- 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`
- Tests organized by module in `tests/test_*.py`, including CLI, YAML, mapping, nested workflows, and discovery
- Timeout tests skip on non-Unix (no `signal.SIGALRM`)
- Use `RecordingHook` pattern to assert event sequences
78 changes: 77 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@

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

For an agent-friendly, runnable YAML example with validation and troubleshooting steps, see the [agent quickstart](docs/agent-quickstart.md).

## Installation

```bash
Expand Down Expand Up @@ -199,6 +201,8 @@ job = Job(workflow=workflow, config=EmptyConfig(), job_configuration=job_config)
result = Runner().run(job)
```

Every declared config field, on root, dependent and mapped tasks alike, must have a value in the `JobConfiguration`; otherwise `Job(...)` raises `WorkflowDefinitionError`.

## Nested Workflows

A workflow can be wrapped as a typed task and used inside a larger workflow. Its input
Expand Down Expand Up @@ -470,6 +474,25 @@ 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`.

Agents can inspect installed task plugins from the CLI without loading every plugin:

```bash
taskmaestro tasks list --json
taskmaestro tasks describe acme.prepare --json
```

`list --json` emits `{"tasks": ["acme.prepare", ...]}` in sorted order; an empty list
means no task plugins are installed. `describe --json` loads only the named plugin
and emits its registered `identifier`, task `name`, `timeout_seconds`, and Pydantic
`input_schema` / `output_schema` (JSON Schema objects). These commands inspect
**installed entry points**, not task classes local to a workflow YAML file. Without
`--json`, `list` prints one identifier per line and `describe` prints indented JSON.
For runtime-only Python objects (such as `ObjectModel[rips.EclipseCase]`), schema
fields include `"not": {}`, `"x-taskmaestro-opaque": true`, and
`"x-taskmaestro-python-type"`. They cannot be supplied as JSON; wire them from
upstream tasks or a Python context instead. Other unsupported schema constructs,
unknown identifiers, and invalid plugins report an error on stderr and exit with status 2.

## YAML Configuration

Workflows can be defined entirely in YAML instead of Python. A `task:` value may be
Expand Down Expand Up @@ -618,9 +641,62 @@ 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
taskmaestro workflow describe workflow.yaml --json
```

`run` prints the final output as JSON and returns a nonzero exit code when the workflow fails. `graph` prints Mermaid markup.
By default, `run` prints the final output as JSON and reports errors on stderr;
`validate` prints a human-readable confirmation. `graph` prints Mermaid markup.

For automation, use `--json` with `validate` or `run`:

```bash
taskmaestro validate workflow.yaml --input input.yaml --json
taskmaestro run workflow.yaml --input input.yaml --json
```

Successful validation emits `{"status":"valid","workflow":"..."}`. A successful
run emits `{"status":"completed","workflow":"...","result":{...}}`. A failure
emits one JSON object on stdout, for example:

```json
{"status":"failed","workflow":"example","failed_task":"prepare","error":{"code":"task_failed","type":"ValueError","message":"Task failed","task":"prepare","field":null,"issues":[]}}
```

Errors contain `code`, exception `type`, a safe `message`, nullable `task` and
`field`, and `issues` (field paths and error codes, without input values).
Loading failures have status `invalid` and `code: "configuration_error"`; task
failures have status `failed` and `code: "task_failed"`. If a completed result
cannot be encoded as JSON (e.g. a Python-only object), `run --json` returns
`code: "serialization_error"`. Missing configuration fields produce `issues`
with code `missing`; when no structured field metadata is available, `field`
is `null`. Use text mode when you need the original exception message.

Exit codes: `0` success, `1` task or result-serialization failure, `2` workflow
configuration failure. In JSON mode logs and ordinary Python `print()` output
from imports/tasks go to stderr, reserving stdout for the result document.
Application code can still write directly to file descriptor 1; JSON mode is
not a sandbox. Error objects intentionally omit raw exception messages, but
application-generated stderr may contain sensitive data.

Inspect a workflow **before its input file is complete** with
`taskmaestro workflow describe workflow.yaml --json`. The JSON contains the
workflow name and result task plus topologically ordered task instances. Each
instance includes its `name`, `python_type`, `depends_on` references (`task` /
`field`), `config_fields`, `required_input_fields`, optional `map` configuration,
and Pydantic input/output schemas. Collection dependencies include their kind
(`positional` or `keyed`) and member references. `python_type` identifies the
loaded class; it may differ from the plugin entry-point identifier used in YAML.

Pass `--input input.yaml` to check that required configuration fields and map
sources are supplied. Then `provided_config_fields` lists **field names only**;
without `--input`, it is `null`. This check does not run tasks or hooks and does
not validate every configured value's runtime type. Loading YAML still imports
Python modules (and nested workflows); **do not inspect untrusted YAML or plugins**
under a privileged account. In JSON mode inspection failures have status
`invalid`, an `error` object, and exit code 2. Without `--json`, inspection
prints indented JSON.

`python -m taskmaestro ...` is equivalent, which is useful when the scripts directory is not on `PATH`.

## Examples

Expand Down
102 changes: 102 additions & 0 deletions docs/agent-quickstart.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
# Agent quickstart

Use this as a short, reproducible path from a typed task to a validated workflow. Commands below run from the **repository root**. The example uses only Taskmaestro's declared dependencies; it needs no API keys or external services.

## Set up

Python 3.12+ is required. In a fresh checkout:

```bash
python3 -m venv .venv
.venv/bin/python -m pip install -e .
```

If you already have Taskmaestro installed in another environment, use that environment's Python instead of `.venv/bin/python` below.

## Inspect the example

- [`examples/agent_quickstart/tasks.py`](../examples/agent_quickstart/tasks.py) defines `AddOne` and `Double`. Each accepts and returns a Pydantic `Number` model. Tasks implement `run(input, ctx)`.
- [`examples/agent_quickstart/workflow.yaml`](../examples/agent_quickstart/workflow.yaml) gives the tasks stable instance names, configures `add_one.value` from the input file, and routes `add_one`'s output to `double`.
- [`examples/agent_quickstart/input.yaml`](../examples/agent_quickstart/input.yaml) supplies `value: 5` to `add_one`.

The YAML loader imports `tasks.AddOne` and `tasks.Double` from the workflow file's directory when invoked through the CLI. No package installation or Python path changes are needed for these local tasks.

## Validate, inspect, run

```bash
.venv/bin/python -m taskmaestro validate examples/agent_quickstart/workflow.yaml --input examples/agent_quickstart/input.yaml
.venv/bin/python -m taskmaestro graph examples/agent_quickstart/workflow.yaml --input examples/agent_quickstart/input.yaml
.venv/bin/python -m taskmaestro run examples/agent_quickstart/workflow.yaml --input examples/agent_quickstart/input.yaml
```

`validate` prints `Workflow 'agent_quickstart' is valid`; `graph` prints Mermaid text with `add_one -->|Number| double`; `run` prints JSON:

```json
{
"value": 12
}
```

The calculation is `(5 + 1) * 2`. The CLI writes errors to stderr and returns a nonzero exit code on failure. `validate` loads and checks the workflow and job configuration **without executing tasks**; `run` executes them.

## Fix a validation error

To reproduce a missing configuration field without changing the checked-in input, create a temporary empty input file:

```bash
printf '{}\n' > /tmp/taskmaestro-agent-bad-input.yaml
.venv/bin/python -m taskmaestro validate examples/agent_quickstart/workflow.yaml --input /tmp/taskmaestro-agent-bad-input.yaml
```

This exits with code 2 and reports:

```text
Configuration error: Job validation failed: Task 'add_one' is missing configuration fields ['value']
```

The workflow declares `value` in `config_fields`, so the input must include `add_one: {value: 5}`. Validate again with the checked-in input, then run. Remove the temporary file when done:

```bash
rm /tmp/taskmaestro-agent-bad-input.yaml
```

## When generating your own workflow

1. Define each task's input and output as Pydantic models in a Python module. Subclass `Task[InputModel, OutputModel]` and implement `run(input, ctx)`.
2. List tasks in a workflow YAML file. Prefer explicit `name:` values; dependency references and top-level input YAML keys refer to **instance names**. Use `depends_on` to connect tasks and `config_fields` for fields supplied from input YAML.
3. Put input data under the task instance name in a separate YAML file. Run `validate`, then `graph`, then `run`. When validation fails, correct the named task/field before retrying.
4. For runtime failures, check stderr for the failing task. A successful validation does not test the task's `run()` logic or guarantee external services are available.

## Discover installed task plugins

If tasks are published by an installed package through the `taskmaestro.tasks` entry-point group, inspect their identifiers and model schemas before generating a workflow:

```bash
.venv/bin/python -m taskmaestro tasks list --json
.venv/bin/python -m taskmaestro tasks describe acme.prepare --json
```

The list command returns a sorted JSON object such as `{"tasks": ["acme.prepare"]}`. The describe command returns the chosen task's `identifier`, `name`, `timeout_seconds`, and Pydantic `input_schema` / `output_schema` JSON Schema objects. Fields containing runtime-only Python objects are marked `x-taskmaestro-opaque` and `x-taskmaestro-python-type`, with `"not": {}` because no JSON value can satisfy them; route these values from upstream tasks rather than inventing JSON input. Replace `acme.prepare` with an identifier from your list; if you have no installed task plugins, the list is empty. The example tasks above are **local Python classes**, not installed plugins, so they will not appear in `tasks list`.

## Inspect a workflow before writing input

Inspection needs only the workflow YAML, so use it to find required fields and routing before creating an input file:

```bash
.venv/bin/python -m taskmaestro workflow describe examples/agent_quickstart/workflow.yaml --json
```

This reports the result task `double`, a root `add_one` with `config_fields: ["value"]`, the dependency on `add_one`, and schemas for each task. Add `--input examples/agent_quickstart/input.yaml` to check that required configuration fields are present; the output includes `provided_config_fields` **names**, not their values. Inspection does not execute tasks or instantiate hooks, but it **does import Python modules** named in YAML (including nested workflows). Only inspect trusted workflow files and plugins. Runtime input values are checked when the tasks run, not completely by inspection.

## Parse results and errors as JSON

Use `--json` with `validate` and `run` when you need a stable response instead of parsing text from stderr:

```bash
.venv/bin/python -m taskmaestro validate examples/agent_quickstart/workflow.yaml --input examples/agent_quickstart/input.yaml --json
.venv/bin/python -m taskmaestro run examples/agent_quickstart/workflow.yaml --input examples/agent_quickstart/input.yaml --json
```

Validation returns `{"status": "valid", "workflow": "agent_quickstart"}`; a successful run returns `{"status": "completed", "workflow": "agent_quickstart", "result": {"value": 12}}`. On failure, stdout contains a single JSON object with status `invalid` (configuration error) or `failed` (execution/serialization error), and an `error` containing `code`, `type`, `message`, `task`, `field`, and `issues`. Missing metadata is `null` or an empty list. Exit codes are 0 for success, 1 for a task/serialization failure, and 2 for an invalid configuration. Error messages omit raw exception details and input values; stderr may still contain application logs or prints. See the [CLI section of the README](../README.md#command-line-interface) for the full contract.

For fan-in, collections, mapping, nested workflows, and the Python builder API, see the [README](../README.md) and the other examples under `examples/`.
2 changes: 2 additions & 0 deletions examples/agent_quickstart/input.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
add_one:
value: 5
19 changes: 19 additions & 0 deletions examples/agent_quickstart/tasks.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
"""Small, dependency-free tasks for the agent quickstart."""

from pydantic import BaseModel

from taskmaestro import ExecutionContext, Task


class Number(BaseModel):
value: int


class AddOne(Task[Number, Number]):
def run(self, input: Number, ctx: ExecutionContext) -> Number:
return Number(value=input.value + 1)


class Double(Task[Number, Number]):
def run(self, input: Number, ctx: ExecutionContext) -> Number:
return Number(value=input.value * 2)
9 changes: 9 additions & 0 deletions examples/agent_quickstart/workflow.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
workflow:
name: agent_quickstart
tasks:
- task: tasks.AddOne
name: add_one
config_fields: [value]
- task: tasks.Double
name: double
depends_on: add_one
3 changes: 2 additions & 1 deletion examples/image_processing/input.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -3,5 +3,6 @@
# Run:
# python examples/image_processing/pipeline.py --yaml --input input.yaml

# Relative paths are resolved against this example's directory.
load_image:
image_path: "taskmaestro.png"
image_path: "../../taskmaestro.png"
29 changes: 20 additions & 9 deletions examples/image_processing/pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,12 +30,14 @@
import hashlib
import struct
from pathlib import Path
from typing import cast

from pydantic import BaseModel

from taskmaestro import (
ExecutionContext,
Job,
JobStatus,
Runner,
Task,
Workflow,
Expand Down Expand Up @@ -209,14 +211,21 @@ def run(self, input: AnalysisInput, ctx: ExecutionContext) -> ImageAnalysis:
# ---------------------------------------------------------------------------


EXAMPLE_DIR = Path(__file__).resolve().parent


class LoadImage(Task[ImageInput, ImagePath]):
"""Resolve the image path to an absolute path."""
"""Resolve the image path to an absolute path.

Relative paths are resolved against this example's directory rather than
the current working directory, so the example runs from anywhere.
"""

name = "load_image"

def run(self, input: ImageInput, ctx: ExecutionContext) -> ImagePath:
ctx.logger.info("Loading image path: %s", input.image_path)
resolved = str(Path(input.image_path).resolve())
resolved = str((EXAMPLE_DIR / input.image_path).resolve())
return ImagePath(path=resolved)


Expand Down Expand Up @@ -280,7 +289,12 @@ def print_report(
outer_workflow: Workflow,
) -> None:
"""Print the analysis report, timings, and Mermaid diagrams."""
report: ReportOutput = result.result # type: ignore[assignment]
if result.status != JobStatus.COMPLETED:
print(f" Job status: {result.status}")
print(f" Failed task: {result.failed_task}")
print(f" Error: {result.error}")
return
report = cast(ReportOutput, result.result)

print("=" * 60)
print(f" {report.title}")
Expand All @@ -303,8 +317,7 @@ def print_report(

def run_python_mode() -> None:
"""Run the pipeline using the Python API."""
_dir = Path(__file__).resolve().parent
image_path = str(_dir / ".." / ".." / "taskmaestro.png")
image_path = "../../taskmaestro.png" # relative to EXAMPLE_DIR

# Build the outer workflow
outer_workflow = (
Expand Down Expand Up @@ -345,20 +358,18 @@ def run_yaml_mode(workflow_path: str, input_path: str) -> None:
def main() -> None:
import argparse

_dir = Path(__file__).resolve().parent

parser = argparse.ArgumentParser(description="Image Processing Pipeline example")
parser.add_argument(
"--yaml",
metavar="FILE",
nargs="?",
const=str(_dir / "workflow.yaml"),
const=str(EXAMPLE_DIR / "workflow.yaml"),
help="Load workflow from a YAML config file (default: workflow.yaml)",
)
parser.add_argument(
"--input",
metavar="FILE",
default=str(_dir / "input.yaml"),
default=str(EXAMPLE_DIR / "input.yaml"),
help="Input YAML file (default: input.yaml)",
)
args = parser.parse_args()
Expand Down
Loading
Loading