baladworkflow
Production-grade Python automation framework. 24 service connectors, 55 batteries-included functions, a resilient workflow engine, signed webhooks and a cron scheduler — with one core dependency.
from baladworkflow import Pipeline, Scheduler
from baladworkflow.tools import GoogleSheets, Slack
from baladworkflow.functions import dedupe_records, filter_records, format_table
sheet, slack = GoogleSheets("<spreadsheet-url>"), Slack()
daily = (
Pipeline("high-value-leads")
.step("read", lambda _: sheet.rows())
.step("dedupe", lambda rows: dedupe_records(rows, key="email"))
.step("filter", lambda rows: filter_records(rows, {"amount__gte": 1000}))
.step("notify", lambda rows: slack.send("#sales", format_table(rows)), on_error="skip")
)
scheduler = Scheduler()
scheduler.cron("0 9 * * mon-fri", daily.run, name="daily-leads")
scheduler.run_forever()
Why this exists
Zapier and n8n are excellent until you need a loop, a real retry policy, a secret that never reaches a log, or a test. Then you write Python — and rewrite the same four things every time: HTTP retries, rate limits, "have I already processed this?", and webhook signature checks.
baladworkflow is those four things, done once and tested, plus the connectors and helpers on top.
| Hand-rolled scripts | Zapier / n8n | baladworkflow | |
|---|---|---|---|
| Version control, diffs, code review | ✅ | ❌ | ✅ |
| Retries with jittered backoff | you write it | ✅ | ✅ |
| Circuit breaker on a failing API | rarely | ❌ | ✅ |
| Durable "already processed" state | you write it | ✅ | ✅ |
| Webhook signature verification | often skipped | ✅ | ✅ |
| Secrets kept out of logs | often not | ✅ | ✅ |
| Dry-run against production creds | ❌ | ❌ | ✅ |
| Per-task pricing | free | per task | free |
Install
pip install baladworkflow
The core needs only httpx. Extras pull in what a particular connector requires:
pip install 'baladworkflow[server]' # FastAPI + uvicorn for the webhook server
pip install 'baladworkflow[google]' # service-account auth (refresh tokens need nothing)
pip install 'baladworkflow[aws]' # S3
pip install 'baladworkflow[sftp]' # SFTP
pip install 'baladworkflow[sql]' # PostgreSQL, MySQL (SQLite works out of the box)
pip install 'baladworkflow[all]' # everything
Start a project:
baladworkflow init my-automation && cd my-automation
# fill in .env, then:
baladworkflow doctor
python automation.py
The four layers
1. Connectors — 24 services
Every connector shares the same constructor contract, retry policy, rate limiter, dry_run mode and health_check(). Learn one, you know all of them.
| Category | Connectors |
|---|---|
Gmail · SmtpMailer · ImapReader |
|
| Spreadsheet / storage | GoogleSheets · GoogleDrive · S3 · Sftp |
| Calendar | GoogleCalendar |
| Messaging | Slack · Discord · Telegram · Twilio |
| Productivity | Notion · Airtable · Trello · Jira |
| Payments / CRM / commerce | Stripe · HubSpot · Shopify |
| Dev | GitHub |
| AI | Claude · OpenAICompatible |
| Data | Database (SQLite/PostgreSQL/MySQL) · RestApi (anything else) |
from baladworkflow.tools import Gmail, Notion, Stripe
gmail = Gmail() # credentials from the environment
for message in gmail.search("is:unread has:attachment", limit=10):
gmail.download_attachments(message, "invoices/")
gmail.mark_read(message["id"])
Notion().create_page("<database-id>", {"Name": "Ada", "Amount": 250, "Tags": ["vip"]})
Stripe().payments(status="succeeded", created_after="last_7_days")
Nested API payloads are flattened for you. Gmail's MIME tree becomes {from, subject, text, html, attachments}; Notion's {"Name": {"title": [{"text": {"content": "Ada"}}]}} becomes {"Name": "Ada"}; Jira, Shopify and Stripe likewise.
2. Functions — 55, all registered and discoverable
baladworkflow functions --search phone
baladworkflow run slugify text='Şirket Raporu 2026' # → sirket-raporu-2026
| Module | # | What it covers |
|---|---|---|
text |
11 | slugify clean_whitespace truncate_text strip_html extract_emails extract_urls extract_phone_numbers render_template similarity_ratio mask_pii format_table |
data |
12 | flatten_dict unflatten_dict deep_merge dig pick_fields rename_fields dedupe_records group_by sort_records filter_records chunk_list diff_records |
validate |
9 | is_valid_email is_valid_url normalize_phone luhn_check is_valid_iban is_valid_tckn validate_schema coerce_types require_fields |
timeutils |
8 | parse_datetime format_datetime humanize_delta relative_window business_days_between add_business_days date_range in_business_hours |
files |
9 | read_csv_rows write_csv to_csv_string load_json_file save_json_file hash_file find_files zip_files human_bytes |
web |
6 | http_request download_file post_webhook parse_rss extract_links build_url |
They exist because each one is a bug people hit in production:
from baladworkflow.functions import coerce_types, dig, filter_records, mask_pii, relative_window
# "1.234,56" and "1,234.56" are both 1234.56 — the CSV import fix
coerce_types({"total": "1.234,56"}, {"total": "float"}) # {'total': 1234.56}
# no more payload.get("a", {}).get("b", {}).get("c")
dig(payload, "items.0.customer.email", default="")
# a filter that can safely come from YAML or an HTTP request
filter_records(rows, {"status": "paid", "amount__gte": 100, "name__contains": "ltd"})
# report windows that keep meaning the right thing next month
start, end = relative_window("last_7_days")
# never log a customer's data
log.info("processing", extra={"body": mask_pii(email_body)})
Every function is introspectable, so the catalogue doubles as LLM tool definitions:
from baladworkflow import tool_schemas
from baladworkflow.tools import Claude
Claude().message(messages, tools=tool_schemas(category="data")) # 12 ready-made tools
3. Engine — pipelines, retries, scheduling, durable state
from baladworkflow import Pipeline, RetryPolicy
flow = (
Pipeline("orders")
.step("fetch", fetch_orders, retry=RetryPolicy(attempts=5))
.filter("unpaid", lambda o: o["status"] == "unpaid")
.map("enrich", add_customer_details)
.branch("route", lambda os: len(os) > 100, if_true=bulk_import, if_false=individual_import)
.step("notify", send_summary, on_error="skip") # a failed Slack post is not fatal
)
result = flow.run()
print(result.ok, result.failed_step, result.to_dict())
Error policies per step: raise (default), skip (pass the input through), stop (end the run), null (continue with None).
Durable state is what makes a scheduled job safe to restart:
from baladworkflow import Deduplicator, get_store
seen = Deduplicator(get_store(), "orders")
for order in fetch_orders():
if seen.is_new(order["id"]): # survives restarts and deploys
process(order)
Scheduling uses a real cron parser — no dependency, aliases and ranges included:
scheduler.cron("*/15 9-17 * * mon-fri", check_queue)
scheduler.cron("@daily", nightly_backup)
scheduler.every(300, poll_inbox, name="inbox")
scheduler.start() # background thread, or .run_forever()
baladworkflow cron '0 9 * * mon-fri' # preview the next firings before you deploy
4. Webhooks — verified, deduplicated, routed
from baladworkflow.webhooks import WebhookRouter, WebhookServer
router = WebhookRouter(
secrets={
"github": os.environ["GITHUB_WEBHOOK_SECRET"],
"stripe": os.environ["STRIPE_WEBHOOK_SECRET"],
}
)
@router.on("stripe", "invoice.payment_failed")
def dunning(event):
slack.send("#billing", f"Payment failed: {event.get('data.object.customer_email')}")
@router.on("github", "push")
def deploy(event):
if event.get("ref") == "refs/heads/main":
trigger_deploy()
WebhookServer(router, port=8000).serve_forever()
What you get without writing it:
- Signature verification for GitHub, Slack, Stripe and Shopify — each provider signs differently, all four are implemented, every comparison is constant-time.
- Replay protection — Slack and Stripe timestamps outside a 5-minute window are rejected; every delivery id is deduplicated durably.
- Handler isolation — an exception in your handler is logged and reported, never returned as a 500 that makes the provider disable your endpoint.
- Wildcards —
@router.on("stripe", "invoice.*").
Already running FastAPI? create_fastapi_app(router) mounts the same routes.
baladworkflow serve --port 8000 --handler myapp:router
Production concerns, handled
Secrets never reach your logs. Credentials are wrapped in Secret, and the JSON log formatter redacts known key names and token patterns (sk-…, xoxb-…, ghp_…, JWTs) anywhere in a record.
>>> print(Secret("hunter2"))
Secret('***')
Dry-run before you trust it. Every mutating request is logged and suppressed:
BALADWORKFLOW_DRY_RUN=true python automation.py
Typed errors, not stack traces from three libraries deep. Everything derives from BaladworkflowError: AuthenticationError, NotFoundError, RateLimitError (carries retry_after), ValidationError, TransportError, RemoteServiceError.
Retries that do not stampede. Exponential backoff with jitter, Retry-After honoured, and a circuit breaker that stops calling a service that is clearly down.
Rate limits respected. Each connector ships the ceiling its provider documents — Slack 1/s, Notion 3/s, Airtable 5/s, Shopify 2/s.
Structured logs. JSON when piped, human-readable in a terminal, correlation ids across a whole run.
CLI
baladworkflow tools -v # connectors, their actions and auth
baladworkflow functions --category data # browse the catalogue
baladworkflow functions slugify # full docs for one function
baladworkflow run normalize_phone phone='(0532) 123 45 67'
baladworkflow doctor # check every configured connector
baladworkflow cron '*/15 9-17 * * mon-fri'
baladworkflow schema --llm -o tools.json # export as LLM tool definitions
baladworkflow serve --handler myapp:router
baladworkflow init # starter .env and example
Configuration
Credentials come from arguments, a .env file, or the environment — in that order.
GOOGLE_CLIENT_ID=... # Gmail, Sheets, Drive, Calendar
GOOGLE_CLIENT_SECRET=...
GOOGLE_REFRESH_TOKEN=...
SLACK_BOT_TOKEN=xoxb-...
NOTION_TOKEN=secret_...
STRIPE_API_KEY=sk_live_...
BALADWORKFLOW_TIMEOUT=30
BALADWORKFLOW_MAX_RETRIES=3
BALADWORKFLOW_DRY_RUN=false
Google needs no client library — the refresh-token grant is a single form POST, implemented directly:
from baladworkflow.tools import authorization_url, exchange_code
print(authorization_url(client_id, ["gmail.send", "sheets"])) # open, approve
tokens = exchange_code(code_from_redirect, client_id, client_secret)
print(tokens["refresh_token"]) # → .env, once, forever
Recipes
Invoice inbox → Sheets → Slack
from baladworkflow import Pipeline
from baladworkflow.tools import Claude, Gmail, GoogleSheets, Slack
from baladworkflow.functions import format_table
gmail, sheet, slack, ai = Gmail(), GoogleSheets("<url>"), Slack(), Claude()
def parse(message):
fields = ai.extract(
message["text"],
{
"vendor": "company that issued the invoice",
"amount": "total amount as a number",
"due_date": "payment due date",
},
)
return {**fields, "email_id": message["id"], "subject": message["subject"]}
(
Pipeline("invoices")
.step("fetch", lambda _: gmail.new_messages("is:unread subject:invoice"))
.map("parse", parse)
.step("store", lambda rows: sheet.append_rows(rows) and rows)
.step("notify", lambda rows: slack.send("#finance", format_table(rows)), on_error="skip")
.run()
)
Nightly CSV → database, with a report on what changed
from baladworkflow.functions import coerce_types, diff_records, read_csv_rows, require_fields
from baladworkflow.tools import Database, SmtpMailer
db = Database("postgresql://user:pass@host/db")
rows = [coerce_types(r, {"amount": "float", "active": "bool"}) for r in read_csv_rows("export.csv")]
valid, invalid = require_fields(rows, ["id", "email"])
changes = diff_records(db.query("SELECT * FROM customers"), valid, key="id")
db.upsert("customers", changes["added"] + changes["changed"], key_columns=["id"])
SmtpMailer(provider="gmail").send(
"ops@acme.com",
"Nightly sync",
f"{len(changes['added'])} new, {len(changes['changed'])} changed, {len(invalid)} rejected",
)
SFTP drop folder → S3 archive
from baladworkflow.tools import S3, Sftp
with Sftp(host="sftp.partner.com", username="acme") as sftp:
for remote in sftp.list_files("/incoming", pattern="*.csv"):
local = sftp.download(remote["path"], "inbox/")["path"]
S3(bucket="archive").upload(local, f"partner/{remote['name']}")
sftp.move(remote["path"], "/processed") # so tomorrow's run skips it
Stale PR nag, every weekday morning
from baladworkflow import Scheduler
from baladworkflow.tools import GitHub, Slack
gh, slack = GitHub(repo="acme/api"), Slack()
def nag():
stale = gh.stale_pull_requests(days=3)
if stale:
slack.send("#eng", "\n".join(f"• <{p['url']}|{p['title']}> — {p['age']}" for p in stale))
scheduler = Scheduler()
scheduler.cron("0 9 * * mon-fri", nag, name="pr-nag")
scheduler.run_forever()
More in examples/.
Testing your automations
The connectors accept an injected transport, so your tests never touch the network:
from baladworkflow.http import HttpClient
from baladworkflow.tools import Slack
def test_alert_posts_to_the_right_channel(transport):
transport.responses.append({"json": {"ok": True, "ts": "1.2", "channel": "C1"}})
slack = Slack(token="x", client=HttpClient(base_url="https://slack.test", transport=transport))
slack.send("#alerts", "disk full")
assert transport.last["json"]["channel"] == "#alerts"
baladworkflow's own suite is 280+ tests, all offline.
Compatibility
Python 3.9 – 3.13, on Linux, macOS and Windows. One required dependency (httpx); everything else is an opt-in extra.
Contributing
git clone https://github.com/EthYusuf/baladworkflow && cd baladworkflow
pip install -e '.[dev]'
pytest && ruff check . && mypy src/baladworkflow
A new connector is one file in src/baladworkflow/tools/: subclass BaseTool, add @register_tool(...), use self.client for requests, and implement health_check(). See docs/adding-a-connector.md.
License
MIT © Yusuf Adin — see LICENSE.
Metadata
Release files for baladworkflow 0.1.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| baladworkflow-0.1.0.tar.gz | 164.8 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| baladworkflow-0.1.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 335.3 kB
Release files / baladworkflow-0.1.0.tar.gz
| Download URL | baladworkflow-0.1.0.tar.gz |
|---|---|
| Size | 164.8 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
2491a7c1def73bebe799c44889c181e7f4956be90248775b2b929b2987c6c097
|
|
BLAKE2b-256 checksum How to use checksums |
5912236bb09361517a1c7b246b62a990c8097858a6f11df8f2ee4ad56324dda9
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.14.6
|
Release files / baladworkflow-0.1.0-py3-none-any.whl
| Download URL | baladworkflow-0.1.0-py3-none-any.whl |
|---|---|
| Size | 170.5 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
8a85face91a3cfab46756f8485517dbb718d0b8779f7da0b03377911720c6ed9
|
|
BLAKE2b-256 checksum How to use checksums |
71cbe52f7144909c3b84bfe7ced8f5a3c07bac28d95e2947770813ef9c8fdc0a
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.14.6
|