Lightweight durable workflow execution for Python backends
Project description
FlowSmith
Lightweight durable workflow execution for Python backends.
FlowSmith fills the gap between fragile custom scripts and heavyweight orchestration systems like Temporal or Airflow. Define backend workflows as code, execute them reliably, and resume them safely after crashes — with zero external infrastructure beyond your existing database.
The problem
Modern backend systems require orchestrating multiple API calls in sequence:
product = call_products_api(product_id)
methods = call_payments_api(product["price"])
payment = create_payment(methods[0]["id"])
notify(user_id, payment["id"])
Writing this is easy. Making it production-safe is not. What happens when:
create_paymentsucceeds but the process crashes beforenotifyruns?call_payments_apitimes out on attempt 1 but would succeed on attempt 2?- You need to know exactly which step failed and why?
Without FlowSmith you write custom retry logic, track state manually, and hope nothing crashes mid-flow. With FlowSmith, every step is persisted, every failure is recoverable, and every execution is observable.
Install
# PostgreSQL
pip install flowsmith[postgres]
# MySQL
pip install flowsmith[mysql]
# Both
pip install flowsmith[postgres,mysql]
Quickstart
1. Configure once at server startup
import flowsmith
flowsmith.configure(
database_url=os.environ["DATABASE_URL"],
pool_min=2, # always 2 connections alive
pool_max=10, # scales to 10 under load
pool_timeout=30, # wait 30s for a free connection before erroring
)
# Optional but recommended — detects crashed nodes automatically
flowsmith.start_watchdog(timeout_seconds=300, interval_seconds=60)
Supports both PostgreSQL and MySQL — FlowSmith detects the right backend from the URL:
DATABASE_URL=postgresql://user:password@host:5432/mydb
DATABASE_URL=mysql://user:password@host:3306/mydb
2. Run migrations
flowsmith migrate
# or explicitly
flowsmith migrate --url postgresql://user:pass@localhost/mydb
This creates two tables in your database: fs_flows and fs_nodes.
3. Define step functions
Steps are plain Python functions. No base classes, no decorators, no magic.
def get_product(ctx):
product = products_api.get(ctx.data["product_id"])
return {
"name": product["name"],
"price": product["price"],
"currency": product["currency"],
}
def get_payment_methods(ctx):
product = ctx.data["get_product"] # output from step 1
methods = payments_api.list(
price=product["price"],
currency=product["currency"],
)
return {"methods": methods}
def create_payment(ctx):
product = ctx.data["get_product"] # output from step 1
methods = ctx.data["get_payment_methods"] # output from step 2
payment = payments_api.create(
method_id=methods["methods"][0]["id"],
amount=product["price"],
)
return {"payment_id": payment["id"]}
def send_notification(ctx):
product = ctx.data["get_product"] # output from step 1
payment = ctx.data["create_payment"] # output from step 3
user_id = ctx.data["user_id"] # original input
email_api.send_receipt(user_id, payment["payment_id"], product["name"])
return {"sent": True}
Each step receives a Context object (ctx). The output of every completed step is available on ctx.data under the step's name. This is how data flows between steps — no argument passing, no globals.
4. Wire up the flow and run it
from flowsmith import Flow, Context
flow = Flow("create_order")
flow.step("get_product", get_product, retries=3)
flow.step("get_payment_methods", get_payment_methods, retries=3)
flow.step("create_payment", create_payment, retries=3)
flow.step("send_notification", send_notification, retries=2)
flow.run(
Context({"product_id": "prod_123", "user_id": "u_456"}),
tracking_id=request.idempotency_key,
)
The Decorator API (v0.4+)
FlowSmith supports a powerful Context-Aware Builder Pattern using decorators. This allows you to define workflows naturally and introduces built-in conditional branching.
from flowsmith.decorators import workflow, step
@workflow("checkout_process")
def process_order():
# Python evaluates this first, so it's Step 1
@step(retries=3)
def fetch_inventory(ctx):
return {"in_stock": True, "price": 100}
# Conditional Branching: This step is skipped if the lambda returns False
@step(condition=lambda ctx: ctx.data["fetch_inventory"]["in_stock"])
def reserve_items(ctx):
pass
# Triggering the nested workflow:
process_order(Context({"order_id": 999}), tracking_id="tr_abc123")
The decorator natively handles retries, pausing on crashes, and skipped steps, exactly like the standard API, but keeps your code clean and composable.
Parallel Execution & Subflows (v0.5+)
FlowSmith supports executing multiple independent steps concurrently utilizing a ThreadPoolExecutor, and formally supports isolated sub-workflows.
from flowsmith.decorators import workflow, step, parallel, subflow
@workflow("analytics_pipeline")
def process_data():
# 1. Parallel execution block
# Everything inside this block runs concurrently!
@parallel
def fetch_phase():
@step
def get_users(ctx): ...
@step
def get_products(ctx): ...
# 2. Sequential step
# Naturally waits for the entire fetch_phase to complete.
@step
def generate_report(ctx): ...
# 3. Trigger a formal sub-workflow
# Safely embeds a child flow process. Blocks parent until child is thoroughly DONE.
@subflow
def alert_email(ctx):
return {
"flow": send_email_flow,
"tracking_id": f"mail_{ctx.data['user_id']}"
}
For power users:
flow.parallel()andflow.subflow()are fully available on the native builder API as well!# 1. Parallel Execution with flow.parallel(): flow.step("get_users", get_users) flow.step("get_products", get_products) flow.step("generate_report", generate_report) # 2. Subflow flow.subflow( "alert_email", send_email_flow, tracking_id=lambda ctx: f"mail_{ctx.data['user_id']}" )
Core guarantees
| Guarantee | What it means |
|---|---|
| Completed steps never re-run | If get_product succeeded, it is skipped on every subsequent retry |
| Resume from last success | A flow that failed on step 3 resumes at step 3, not step 1 |
| At-least-once execution | If a process crashes mid-step, that step will re-run on resume |
| Full execution trace | Every step's input, output, error, and attempt count is stored |
| Idempotent trigger | Calling flow.run() with the same tracking_id resumes, never duplicates |
Note on at-least-once: FlowSmith guarantees completed steps are never re-run. Steps that were in progress when a crash happened will be retried. Make your step functions idempotent (safe to run more than once) for full crash safety.
The watchdog
The watchdog solves a subtle but critical problem: what happens when a process crashes after a node is marked RUNNING but before it is marked COMPLETED?
Without intervention, that node stays RUNNING in the database forever — the flow can never be resumed because FlowSmith sees it as still in progress.
The watchdog is a background thread that periodically scans for nodes that have been RUNNING longer than your configured timeout. When it finds one, it marks the node FAILED and the parent flow FAILED, so the next flow.run() call can resume correctly.
10:00:02 Node starts → status = RUNNING
10:00:02 Server crashes → node frozen in RUNNING
...5 minutes...
10:05:00 Watchdog wakes → finds node RUNNING for > 5 min
10:05:00 Watchdog marks → node = FAILED, flow = FAILED
10:05:30 Server restarts → flow.run(tracking_id="same-id")
10:05:30 FlowSmith sees → step 1 COMPLETED (skip)
10:05:30 FlowSmith sees → step 2 FAILED (retry from here) ✓
Start it at server startup:
flowsmith.configure(database_url=os.environ["DATABASE_URL"])
flowsmith.start_watchdog(
timeout_seconds=300, # how long before a RUNNING node is considered stuck
interval_seconds=60, # how often to scan
)
Set timeout_seconds to at least 2x your slowest step's expected runtime. If your payment API can legitimately take 2 minutes, use a timeout of at least 5 minutes.
The watchdog runs as a daemon thread — it stops automatically when your process exits. No cleanup required.
Resume behaviour
FlowSmith uses tracking_id to identify a flow execution. Pass the same tracking_id to resume:
# First run — crashes on step 3
flow.run(ctx, tracking_id="order_abc123")
# Resume — step 1 and 2 skipped automatically, resumes at step 3
flow.run(ctx, tracking_id="order_abc123")
| Flow status on second call | Behaviour |
|---|---|
| Not found | Fresh execution |
RUNNING or FAILED |
Resume from last completed step |
COMPLETED |
Raises FlowAlreadyCompleted |
Database backends
FlowSmith supports PostgreSQL and MySQL via the same StorageBackend interface. The backend is selected automatically from your database URL.
| URL prefix | Backend |
|---|---|
postgresql:// or postgres:// |
PostgresStorage |
mysql:// |
MySQLStorage |
For testing, use InMemoryStorage — no database required:
from flowsmith.storage import InMemoryStorage
flow = Flow("test_flow", storage=InMemoryStorage())
flow.step("my_step", my_fn)
flow.run(ctx, tracking_id="test-1")
Local development
Note:
- Requires Docker for integration tests
- Databases must be running before migrations/tests
- Uses PostgreSQL + MySQL via docker-compose
- CLI works on Windows, macOS, and Linux
# Start both databases
make db-up
# Run migrations
make migrate-postgres
make migrate-mysql
# Install dev dependencies
make install
# Unit tests — no database needed, runs in milliseconds
make test-unit
# Integration tests — requires running databases
make test-integration
# All tests with coverage
make test
# Smoke tests — execute realistic end-to-end execution flow
python smoke_test.py
# Stop both databases
make db-down
Project structure
flowsmith/
├── flowsmith/
│ ├── __init__.py # public API: configure(), start_watchdog(), Flow, Context
│ ├── config.py # global configure(), start_watchdog(), get_storage()
│ ├── flow.py # Flow class — step registration and run()
│ ├── executor.py # step execution loop, retry, resume logic
│ ├── watchdog.py # background thread for stuck-node detection
│ ├── context.py # Context — shared state between steps
│ ├── step.py # Step dataclass
│ ├── exceptions.py # FlowSmithNotConfigured, StepFailed, FlowAlreadyCompleted
│ ├── storage/
│ │ ├── base.py # StorageBackend ABC — implement this to add a new backend
│ │ ├── postgres.py # PostgreSQL backend
│ │ ├── mysql.py # MySQL backend
│ │ └── memory.py # In-memory backend for tests
│ ├── models/
│ │ ├── flow_record.py # FlowRecord dataclass
│ │ └── node_record.py # NodeRecord dataclass
│ ├── migrations/
│ │ ├── postgres/ # PostgreSQL migration files
│ │ └── mysql/ # MySQL migration files
│ └── contrib/ # Optional integrations
└── tests/
├── unit/ # Fast tests using InMemoryStorage — no infrastructure needed
└── integration/ # Real database tests — requires docker-compose
Roadmap
| Version | Scope |
|---|---|
| v0.1.0 | Core engine, InMemoryStorage, sequential execution |
| v0.2.0 | PostgresStorage, MySQLStorage, migrate CLI, watchdog |
| v0.3.0 | Connection pooling, retry backoff strategies, per-step timeout |
| v0.3.1 | Bugfixes: stuck-node SQL, cross-platform thread handling, CI on PRs, py.typed |
| v0.4.0 | Decorator API builder pattern, conditional branching |
| v0.5.0 | Parallel step execution, First-class blocking sub-flows |
| v0.6.0 | PyPI Public Release, Community outreach, Strict Typing (mypy) ← current |
| v1.0.0 | Stable API, core feature completion, long-term stability |
License
MIT
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 flowsmith-0.6.0.tar.gz.
File metadata
- Download URL: flowsmith-0.6.0.tar.gz
- Upload date:
- Size: 28.3 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.3
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
7d0df3dbcbc32fd8bd5c1d6360a2824a2f4eacb5a75c1fef50554774a3a4cd3c
|
|
| MD5 |
87656e79986548ee7b0280db2e2896d2
|
|
| BLAKE2b-256 |
5835305a79d63158ed1e1eaa5167b0cad2ab83dfe1403c0ee4fe1e88e1513819
|
File details
Details for the file flowsmith-0.6.0-py3-none-any.whl.
File metadata
- Download URL: flowsmith-0.6.0-py3-none-any.whl
- Upload date:
- Size: 31.4 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.3
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
3acd48fa93dff263cdd1fc586682c5b75bbddccc089753a844c894c5ddf6e461
|
|
| MD5 |
d489cf17ed73b4cdbb93ff996c89843a
|
|
| BLAKE2b-256 |
840595dbfb9636678ef847619666e23da6c65ba78e6b7ed8d21a10753704568e
|