A lightweight, developer-first workflow & task orchestration engine for Python.
Project description
auto-workflow
A lightweight, zero-bloat, developer‑first workflow & task orchestration engine for Python.
Status: Alpha (APIs stabilizing). Goal: Production‑grade minimal core with pluggable power features.
Quick links
- Docs (GitHub Pages): https://stoiandl.github.io/auto-workflow/
- Repository: https://github.com/stoiandl/auto-workflow
Table of Contents
- Why Another Orchestrator?
- Philosophy & Design Principles
- Feature Overview
- Quick Start
- Core Concepts
- Execution Modes
- Building Flows & DAGs
- Dynamic Fan‑Out / Conditional Branching
- Result Handling, Caching & Idempotency
- Retries, Timeouts & Failure Semantics
- Hooks, Events & Middleware
- Configuration & Environment
- Observability (Logging, Metrics, Tracing)
- Extensibility Roadmap
- Security & Isolation Considerations
- Comparison with Airflow & Prefect
- Project Structure (Proposed)
- FAQ
- Roadmap
- Contributing
- Versioning & Stability
- License
- Examples Overview
Why Another Orchestrator?
Existing platforms (Airflow, Prefect, Dagster, Luigi) solve orchestration at scale—but often at the cost of:
auto-workflow targets a different sweet spot:
Be the simplest way to express, execute and evolve complex, dynamic task graphs directly in Python—locally first—while remaining extensible to production constraints.
Core values: No mandatory DB, no daemon, no CLI bureaucracy, opt‑in persistence, first-class async, predictable concurrency, explicit data flow.
Philosophy & Design Principles
| Principle | Description |
|---|---|
| Minimal Core | Ship only primitives: Task, Flow (DAG), Executor, Runtime. Everything else is a plugin or optional layer. |
| Python Native | Flows are plain Python; no YAML DSL or templating required. |
| Deterministic by Default | Task graph shape should be reproducible given the same inputs. Explicit APIs for dynamic fan‑out. |
| Composable | Tasks are small units; flows can nest; subgraphs are reusable. |
| Extensible | Storage adapters, retry policies, and event sinks are pluggable via clean interfaces. |
| Progressive Adoption | Use it as a simple task runner first; layer complexity only when needed. |
| Observability Hooks | Logging / metrics / tracing surfaces are unified and optional. |
| Zero Hidden State | No implicit global registry; registration is explicit or decorator‑driven with clear import semantics. |
| Performance Conscious | Support high‑throughput local pipelines via async & thread/process pools with minimal overhead. |
Feature Overview
Core capabilities:
- Task Definition:
@taskdecorator with retry, timeout, caching, and execution mode options - Flow Orchestration:
@flowdecorator for building DAGs with automatic dependency resolution - Dynamic Fan-Out:
fan_out()for runtime task creation based on upstream results - Multiple Execution Modes: async, thread pool, and process pool execution
- Caching & Artifacts: Task result caching and large result persistence
- Observability: Built-in logging, metrics, tracing, and event system
- Configuration: Environment-based config with structured logging
- CLI Tools: Run, describe, and list flows via command line
- Secrets Management: Pluggable secrets providers
- Failure Handling: Configurable retry policies and failure propagation
Quick Start
Install from PyPI:
pip install auto-workflow
Or for local development with Poetry:
Run tests locally:
poetry run pytest --cov=auto_workflow --cov-report=term-missing
Define Tasks
from auto_workflow import task, flow, fan_out
@task
def load_numbers() -> list[int]:
return [1, 2, 3, 4]
@task
def square(x: int) -> int:
return x * x
@task
def aggregate(values: list[int]) -> int:
return sum(values)
@flow
def pipeline():
nums = load_numbers()
# Dynamic fan-out: create tasks for each number
squared = fan_out(square, nums)
return aggregate(squared)
if __name__ == "__main__":
result = pipeline.run()
print(result)
CLI
python -m auto_workflow run path.to.pipeline:pipeline
List and describe flows:
python -m auto_workflow list path.to.module
python -m auto_workflow describe path.to.pipeline:pipeline
If installed via pip, a console script is also available:
auto-workflow run path.to.pipeline:pipeline
Core Concepts
Task
Unit of work: a pure (or side-effecting) Python callable that declares inputs & returns outputs. Decorated with @task for metadata: name, retries, timeout, tags, cache key fn.
Flow (or Pipeline)
Container for a DAG of tasks. May be defined with @flow decorator wrapping a function whose body builds the dependency graph during invocation. Supports nested flows.
DAG
Directed acyclic graph where edges represent data or control dependencies. Construction is implicit via using outputs of tasks as inputs to other tasks (like Prefect) but without global mutable state.
Execution Context
Available via from auto_workflow.context import get_context() inside a task for run metadata, logger, parameters.
Parameters
Flow-level runtime parameters passed at .run(params={...}) enabling configurability without environment variables.
Artifacts
Structured results (maybe large) that can optionally be stored externally; default is in‑memory pass‑through.
Execution Modes
Tasks run using one of three simple modes:
Mode selection:
Building Flows & DAGs
Flows are defined using the @flow decorator:
@flow
def my_flow():
a = task_a()
b = task_b(a)
c = task_c(a, b)
return c
Task dependencies are determined automatically by passing task invocation results as arguments to other tasks.
Dynamic Fan-Out / Conditional Branching
Dynamic forks are explicit to preserve introspection & safety:
@task
def split_batches(data: list[int]) -> list[list[int]]: ...
@task
def process_batch(batch: list[int]) -> int: ...
@task
def combine(results: list[int]) -> int: return sum(results)
@flow
def batch_flow(data: list[int]):
batches = split_batches(data)
results = fan_out(process_batch, iterable=batches, max_concurrency=8)
return combine(results)
@flow
def conditional_flow(flag: bool):
a = task_a()
if flag:
b = task_b(a)
else:
b = task_c(a)
return b
fan_out constructs a bounded dynamic subgraph; a future fan_in utility may allow explicit barrier semantics.
Result Handling, Caching & Idempotency
Tasks support caching with TTL and artifact persistence for large results:
@task(cache_ttl=3600) # Cache for 1 hour
def expensive(x: int) -> int:
return do_work(x)
@task(persist=True) # Store large results via artifact store
def produce_large_dataset() -> dict:
return {"data": list(range(1000000))}
Retries, Timeouts & Failure Semantics
Per-task configuration:
@task(retries=3, retry_backoff=2.0, retry_jitter=0.3, timeout=30)
def flaky(): ...
Failure policy options:
FAIL_FAST: Stop on first error (default)CONTINUE: Continue executing independent tasks
Hooks, Events & Middleware
Lifecycle hook points:
Middleware chain (similar to ASGI / HTTP frameworks):
def timing_middleware(next_call):
async def wrapper(task_ctx):
start = monotonic()
try:
return await next_call(task_ctx)
finally:
duration = monotonic() - start
task_ctx.logger.debug("task.duration", extra={"ms": duration*1000})
return wrapper
Event bus: structured events with pluggable subscribers for custom logging and monitoring.
Configuration & Environment
Configuration via environment variables:
AUTO_WORKFLOW_LOG_LEVEL: Set logging level (default: INFO)AUTO_WORKFLOW_DISABLE_STRUCTURED_LOGS: Disable structured loggingAUTO_WORKFLOW_MAX_DYNAMIC_TASKS: Limit dynamic task expansion
See docs/configuration.md for full details.
Connectors
Production-grade connectors are available as optional extras. Today we ship Postgres (psycopg3 pool-backed) and ADLS2 (Azure Data Lake Storage Gen2). SQLAlchemy helpers for Postgres are optional.
Install:
- pip
pip install "auto-workflow[connectors-postgres]"
# Optional: for SQLAlchemy helpers (engine/session/sessionmaker)
pip install "auto-workflow[connectors-sqlalchemy]"
- Poetry
poetry add auto-workflow -E connectors-postgres
# Optional: for SQLAlchemy helpers
poetry add auto-workflow -E connectors-sqlalchemy
# ADLS2 (Azure SDK)
poetry add auto-workflow -E connectors-adls2
# Or install all available connector extras
poetry add auto-workflow -E connectors-all
Usage (Postgres):
from auto_workflow.connectors import postgres
with postgres.client("default") as db:
rows = db.query("SELECT 1 AS x")
print(rows) # [{"x": 1}]
SQLAlchemy (optional extra):
from auto_workflow.connectors import postgres
db = postgres.client("default")
engine = db.sqlalchemy_engine()
from sqlalchemy.orm import Session
with Session(engine) as s:
# ... use ORM as usual ...
pass
See the Connectors page in the docs for more examples, including streaming iteration and reflection: https://stoiandl.github.io/auto-workflow/connectors/
For local testing with Postgres via Docker Compose and running the full suite, see docs/testing.md.
ADLS2 (Azure Data Lake Storage Gen2)
Install:
pip install "auto-workflow[connectors-adls2]"
# or for Poetry
poetry add auto-workflow -E connectors-adls2
Quick example:
from auto_workflow.connectors import adls2
with adls2.client("default") as fs:
fs.make_dirs("bronze", "incoming/demo", exist_ok=True)
fs.upload_bytes(
container="bronze",
path="incoming/demo/sample.csv",
data=b"id,name\n1,alice\n2,bob\n",
content_type="text/csv",
)
data = fs.download_bytes("bronze", "incoming/demo/sample.csv")
print(data.decode())
Environment config (DEFAULT profile):
# Option A: connection string
export AUTO_WORKFLOW_CONNECTORS_ADLS2_DEFAULT__CONNECTION_STRING="DefaultEndpointsProtocol=..."
# Option B: account_url + DefaultAzureCredential
export AUTO_WORKFLOW_CONNECTORS_ADLS2_DEFAULT__ACCOUNT_URL="https://<acct>.dfs.core.windows.net"
export AUTO_WORKFLOW_CONNECTORS_ADLS2_DEFAULT__USE_DEFAULT_CREDENTIALS=true
End-to-end example flow: examples/adls_csv_flow.py.
Observability (Logging, Metrics, Tracing)
Built-in observability features:
- Structured Logging: Automatic JSON-formatted logging with task/flow context
- Metrics: Pluggable metrics providers (in-memory and custom backends)
- Tracing: Task and flow execution spans for performance monitoring
- Events: Pub/sub event system for task lifecycle hooks
- Middleware: Chain custom logic around task execution
Extensibility
| Extension | Interface | Status |
|---|---|---|
| Storage backend | ArtifactStore |
✅ Implemented |
| Cache backend | ResultCache |
✅ Implemented |
| Metrics provider | MetricsProvider |
✅ Implemented |
| Tracing adapter | Tracer |
✅ Implemented |
| Secrets provider | SecretsProvider |
✅ Implemented |
| Event middleware | Middleware chain | ✅ Implemented |
| Executor plugins | BaseExecutor |
Future |
| Scheduling layer | External module | Future |
| UI / API | Optional service | Future |
Security & Isolation Considerations
Comparison with Airflow & Prefect
| Aspect | auto-workflow | Airflow | Prefect |
|---|---|---|---|
| Requires DB / Scheduler | No (local in-process) | Yes | No (cloud optional) |
| First-class async | Yes (core) | Limited | Yes |
| Dynamic DAG at runtime | Explicit fan-out | Limited / brittle | Supported |
| Footprint | Minimal deps | Heavy | Moderate |
| UI bundled | No (optional) | Yes | Yes |
| Plugin surface | Lean, Pythonic | Large | Large |
| Setup time | Seconds | Minutes+ | Minutes |
Project Structure (Proposed)
auto_workflow/
__init__.py
tasks.py # @task decorator & Task definition
flow.py # Flow abstraction & @flow decorator
dag.py # Internal DAG model
execution.py
base.py
async_executor.py
thread_executor.py
runtime/
context.py
scheduler.py # Lightweight topological / async scheduler
middleware/
events/
caching/
storage/
observability/
tests/
examples/
docs/
FAQ
Q: Is persistence required? No—default run is ephemeral in memory.
Q: Can I dynamically create thousands of tasks? Yes, but bounded; guardrails (max-dynamic-tasks) will protect runaway expansion.
Q: How are circular dependencies prevented? DAG builder performs cycle detection before execution.
Q: Do I need decorators? No; you can manually wrap callables into Tasks if you prefer pure functional style.
Q: How does it serialize arguments across processes? Uses cloudpickle for process execution mode.
Q: Scheduling / cron? Out of core scope—provide a thin adapter so external schedulers (cron, systemd timers, GitHub Actions) can invoke flows.
Roadmap
Examples Overview
Explore runnable examples in examples/ (also rendered in the online docs):
| File | Concept Highlights |
|---|---|
data_pipeline.py |
Basic ETL flow with persistence & simple mapping |
concurrent_priority.py |
Priority scheduling & mixed async timings |
dynamic_fanout.py |
Runtime fan-out expansion and aggregation |
retries_timeouts.py |
Retry + timeout interplay demonstration |
secrets_and_artifacts.py |
Secrets provider & artifact persistence usage |
tracing_custom.py |
Custom tracer capturing spans & durations |
Run any example:
python examples/tracing_custom.py
For more narrative documentation see the Examples page.
Contributing
Contributions are welcome once the core API draft solidifies. Until then:
- Open an issue to discuss proposals.
- Keep changes atomic & well-tested.
- Adhere to Ruff formatting & lint rules (pre-commit enforced).
- Add or update examples & docs for new features.
See CONTRIBUTING.md for detailed contribution guidelines.
Versioning & Stability
Pre-1.0: Breaking changes can occur in minor releases. After 1.0 we will follow Semantic Versioning.
Migration notes will be maintained in CHANGELOG.md (to be added).
License
This project is licensed under the GNU General Public License v3.0. See LICENSE for details. If you need alternate licensing, open an issue to discuss.
Legal / Disclaimer
This project is experimental; do not deploy to critical production paths until a 1.0 release is tagged. Feedback welcomed to refine design decisions.
Happy orchestrating! 🚀
Note on releases: Publishing to PyPI is performed by CI when a tag vX.Y.Z is pushed to the repository. We use PyPI Trusted Publishing; no local uploads are required.
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file auto_workflow-0.1.2.tar.gz.
File metadata
- Download URL: auto_workflow-0.1.2.tar.gz
- Upload date:
- Size: 192.1 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
f13f5b5f52cd135d1315f2dffc68bb5d1a2b00046a9b7ae71f595638f189af6f
|
|
| MD5 |
2adef7545241f690adf985340ed26058
|
|
| BLAKE2b-256 |
b4740ca0d2489eefcd0a2b4d6ee0183cf22372b86127cfe8e2ba3aed3e29b16b
|
File details
Details for the file auto_workflow-0.1.2-py3-none-any.whl.
File metadata
- Download URL: auto_workflow-0.1.2-py3-none-any.whl
- Upload date:
- Size: 67.7 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
b0b3f554f37bf2a4687ebe8bf7f6f593804a2e256baa9b6624c7081494dcc550
|
|
| MD5 |
b3207a77a2ff331d0b681ce3b0d10f58
|
|
| BLAKE2b-256 |
b1409659a77cc8d942ab3947a7945f4e9ddd70a9bbabf8d9273eed161660259f
|