Skip to main content
Pre-release

This release is a pre-release and may not be stable for production use.

Buelon

Buelon logo

Buelon is a Python orchestration system with a small scripting language (a DML) for managing large amounts of I/O-heavy work — API calls for ETL and ELT, and other programs that need coordinated Python and/or SQL execution.

A hub holds the job queue. Workers connect to it, pull jobs, run them, and hand the results back. Jobs form a DAG: a job's return value becomes its children's arguments.

Table of Contents

Installation

pip install buelon

That's it. This installs the bue CLI (boo and pete are aliases). Check the install with bue --version.

Python 3.10 or newer is required.

Quick Start

Everything below runs in one directory. Each command reads its configuration from .boo/settings.yaml in the current working directory.

# 1. Create .boo/settings.yaml
bue init

# 2. (optional) edit .boo/settings.yaml -- host, port, scopes
$EDITOR .boo/settings.yaml

# 3. Start the hub. It holds the queue; leave it running.
bue hub

# 4. In another terminal, start a worker (start as many as you like)
bue worker

# 5. Upload a pipeline
bue upload -f example.bue

# 6. Watch it
bue status            # one-shot
bue status -s         # refresh every 3 seconds
bue web               # web UI on http://localhost:11011

A finished pipeline disappears from bue status: once every job in a DAG has succeeded, the hub drops the whole DAG, so total: 0 means "everything completed", not "nothing was uploaded".

See Learn by Example for an example.bue / example.py pair that runs as written.

Architecture

There are two long-running processes, and both keep all their state in memory:

process command what it does
hub bue hub Holds the job queue, the job graph, and every job's result. One per cluster.
worker bue worker Connects to the hub, pulls jobs, runs them, reports back. Any number.

Every other command is a short-lived client that connects to the hub: upload, status, errors, reset, delete, run-job, and the web UI.

Job results are held in the hub's memory and passed to child jobs from there. There is no separate results store — bue bucket still exists as a standalone key/value server, but nothing on the hub/worker path talks to it.

Hub state is snapshotted to disk. By default the hub writes .auto_save/snapshot every ten minutes, plus once more on a clean shutdown (including SIGTERM, so docker stop and systemctl stop are safe), and reloads it on startup. A crash loses at most one interval. BUELON_AUTO_SAVE=false turns both halves off; BUELON_AUTO_SAVE=load-only restores a snapshot without writing one back.

Scopes and priority

Every job has a scope (a free-form name) and a priority (any integer; 0-100 is the usual range, but nothing is clamped and negatives are fine). A worker only pulls jobs whose scope is in its own scopes list, highest priority first. That is how you keep heavy jobs on big machines, or stop one misbehaving pipeline from starving everything else.

Configuration

Configuration lives in .boo/settings.yaml, relative to the directory each command is run from. bue init writes one with the defaults; bue where prints the path it will use.

The hub/worker path reads nothing else. In particular there is no .env support for hub/worker configuration and no command-line flags for host or port — the CLI parses with parse_known_args(), so bue hub -b 0.0.0.0:65432 is accepted and then silently ignored. Edit the yaml.

hub:
  host: 0.0.0.0        # interface the hub binds
  port: 65432
  encryption: faster   # faster | secure | off -- MUST match on the hub and every worker

worker:
  host: localhost      # the hub's address, as seen by this machine
  port: 65432
  scopes: production-heavy,production-small,default   # comma-separated, no spaces
  reverse: false       # pull the lowest-priority scope first instead of the highest
  info:
    name: Worker       # shown in `bue web`'s worker list

bucket:                # only used by `bue bucket`, which nothing else talks to
  server: {use: true, path: .boo/bucket, host: 0.0.0.0, port: 61535}
  client: {use: true, host: localhost, port: 61535}
  postgres: {use: false, table: buelon_bucket, persistent_path: __PERSISTENT__}

postgres:              # NOT read by anything -- see Known Defects
  host: localhost
  port: 5432
  username: XXXXX
  password: XXXXX
  database: XXXXX

hub.port and worker.port must match, and every client command (upload, status, …) uses the worker block to find the hub.

hub.encryption

The wire format for every hub/worker connection. It lives under hub: but is read by both ends — the hub and every worker, CLI command and bue web process — so there is one value to keep in step rather than two that can disagree.

value wire format
faster (default) AES-GCM.
secure AES-GCM, then bz2 over the ciphertext. bisocket's own default.
off Plaintext. Only on a network you fully trust.

secure is slower and bigger than faster here, not safer: the bz2 pass runs after encryption, so it compresses ciphertext, which is incompressible — a 25-job batch measures ~0.65 ms/frame of hub CPU and comes out ~26% larger on the wire. Buelon already bz2- compresses job batches itself before they reach the transport, which is why the second pass buys nothing. Both modes use the same AES-GCM encryption and the same CRYPTO_KEY.

The hub and every worker must agree. A mismatch is refused with EncryptionMismatch naming both modes — it is not negotiated — so upgrading from a version that predates this setting (or changing the value) means restarting the hub and all its workers together, not a rolling restart. To upgrade without a coordinated restart, set encryption: secure everywhere first, then switch to faster when you can take the hub down.

Leave the value empty (encryption:) to defer to $BISOCKET_ENCRYPTION, and to bisocket's secure default if that is unset too.

Environment variables

variable default effect
BISOCKET_ENCRYPTION — Wire format, consulted only when hub.encryption in the yaml is left empty. Same values as that setting.
CRYPTO_KEY an insecure built-in default Transport encryption key. Set this in production. Every hub, worker and CLI invocation must use the same value; a mismatch fails the connection with EncryptionMismatch.
BUELON_SETTINGS_PATH .boo/settings.yaml Full path to the settings file.
BUELON_DIR_PATH .boo Directory every state file lives under — settings, the parser's scratch files, the bucket store. Setting it also disables the .bue/ → .boo/ migration below.
BUELON_AUTO_MIGRATE true Set to false to skip the one-time .bue/ → .boo/ copy described below.
BUELON_AUTO_SAVE true Set to false to disable hub snapshots entirely — nothing is written and nothing is restored on startup. load-only restores an existing snapshot but never overwrites it.
BUELON_AUTO_SAVE_PATH .auto_save Directory the hub snapshot is written to.
BUELON_AUTO_SAVE_INTERVAL 600 Seconds between snapshots.
BUELON_RETRY_BACKOFF_BASE 5 Seconds before a job's first retry. Each further attempt doubles it. 0 retries immediately. Read by the hub.
BUELON_RETRY_BACKOFF_MAX 300 Ceiling on that doubling, in seconds.
BUELON_HANDBACK_DELAY 5 Seconds a job that returns pending waits before it is offered again. Constant, not doubling — a poll is not a failure. 0 re-queues immediately. Read by the hub and by bue run -f.
BOO_WEB_HOST / BOO_WEB_PORT localhost / 11011 Where bue web listens.
POSTGRES_HOST, POSTGRES_PORT, POSTGRES_USER, POSTGRES_PASSWORD, POSTGRES_DATABASE localhost/5432/… Connection used by postgres jobs (see Supported Languages). Read by workers, not by the hub.
ENV_PATH .env A .env file at this path is loaded, if python-dotenv is installed. Only the variables in this table have any effect.

The state directory is created in the current working directory the first time something needs it — the parser uses it for scratch files, and bue upload / bue run will make it.

It used to be called .bue/. buelon is named after Buelon Rexford Moss, whose nickname was boo; bue was a misspelling. On startup, if a non-empty .bue/ is present and .boo/ is not, buelon copies the one to the other and prints a line saying so. .bue/ is left where it is — so a downgrade still finds its state — but it stops being read from that moment, so the two drift apart. Delete it once you are happy. Set BUELON_AUTO_MIGRATE=false, or point BUELON_DIR_PATH somewhere explicit, to skip this.

Command Reference

bue init                    create .boo/settings.yaml
bue where                   print the settings.yaml path in use

bue hub                     run the hub (foreground)
bue worker                  run a worker (foreground)
bue work                    same as `bue worker`
bue run-job -j <job_id>     run one job by id, once, then exit

bue upload -f <file.bue>    build the pipeline here, send its jobs to the hub
bue submit -f <file.bue>    send the script instead, and build it on a worker
                            (-s SCOPE picks which; default: last of worker.scopes)
bue run -f <file.bue>       run a pipeline locally, start to finish, with no hub

bue status                  one-shot job counts
bue status -s               refresh every 3 seconds
bue status -s -l            ...and print a permanent line every 15 minutes
bue errors                  print every errored job with its traceback
bue reset                   move errored jobs back onto the queue
bue delete                  cancel jobs belonging to errored pipelines
bue delete --all            delete every job on the hub (prompts; -y to skip)
bue web [-o]                web UI on :11011 (-o opens a browser)

bue bucket                  run the standalone bucket server
bue example                 write example.bue / example.py / demo.py into the cwd
bue demo                    run a bucket, a hub and three workers in one process
bue joke                    a boo joke
bue --version               print the version

bue repair and bue test are accepted and do nothing.

bue worker and bue work exit on their own after 20 minutes. That is deliberate (a periodic restart drops any leaked memory or module state), and it is not configurable. Run them under a supervisor that restarts them — systemd with Restart=always, a Docker restart policy, or a shell loop.

Supported Languages

A job's language is the second line of its definition. Three are available:

  • python (also python3, py)
  • sqlite3 (also sqlite)
  • postgres (also postgresql, pg)

For python, the third line is the function to call. For the SQL languages it is the table name the incoming rows are loaded under, and the job's code is a query against that table; the query's result set becomes the job's return value.

postgres here means "run this job's SQL on Postgres", not "store Buelon's state in Postgres". There is no Postgres backend for the hub. Postgres jobs connect using the POSTGRES_* environment variables on the worker that runs them, and require psycopg2-binary and asyncpg.

Job return values

Whatever a Python job returns is sent to the hub and handed to its children, so it has to survive JSON serialization: dicts, lists, strings, numbers, booleans, None. Return an object that cannot be — a uuid.UUID, a set, a datetime — and that job fails with

job 'request' (a699…) returned a value that cannot be sent to the hub:
Object of type UUID is not JSON serializable

visible in bue errors. Only that job fails; the rest of its batch is unaffected. Convert to a primitive (f'{uuid.uuid4()}', dt.isoformat(), list(s)) before returning.

A job can also return a Result to control what the hub does next:

from buelon.core.step import Result, StepStatus

return Result(status=StepStatus.pending)   # not ready; re-queue me and try again later
return Result(status=StepStatus.reset)     # start this chain over from its root
return Result(status=StepStatus.cancel)    # drop this chain

pending is the important one: it is how you poll a slow API without holding a worker slot for the whole wait. The hub holds the job back for BUELON_HANDBACK_DELAY seconds (5 by default) before offering it again, so the poll is a poll rather than a hot loop, and counts the hand-backs — bue status reports them as handed back, and the web UI as Handed Back. There is no limit unless you set !max_handbacks. bue run -f applies the same delay and the same ceiling, so a polling pipeline behaves locally the way it will on the cluster; it steps over a waiting job and runs the rest of the pipeline meanwhile.

Those three are the whole list. The other StepStatus members — queued, working, success, error, unknown — are hub bookkeeping, not things a job returns. queued in particular is not a slower pending: it means "this job is blocked on a parent that has not finished yet", and the hub sets and clears it as parents complete. Returning it from a job is recorded as an error, because by the time a worker holds a job its parents are already done, so there is nothing left to wait for. Use pending to be tried again.

Learn by Example

The two files below are the ones used to verify this README. Write both into the same directory, then bue upload -f example.bue. (bue example writes a second, slightly fuller working example into the current directory — same pipeline shape, plus a sqlite3 job and a .boo/settings.yaml.)

example.bue

# Defaults for every job in this file.
!scope default
!timeout 20 * 60

# A job definition: name, language, function/table name, then the code
# (a file path, or inline code between backticks).
accounts:
    python
    accounts
    example.py

# Or define several at once out of the same file.
import python (
    request_report as request,
    get_status as status,
    get_report
        as download
        !priority 9,
    upload_to_db as upload
) example.py

# SQL jobs take their input table under the name you give here.
manipulate_data:
    sqlite3
    some_table
    `
SELECT
    *,
    CASE WHEN sales = 0 THEN 0.0 ELSE spend / sales END AS acos
FROM some_table
`

# Pipes say what order jobs run in, and pass each job's return value to the next.
accounts_pipe = | accounts
api_pipe = request | status | download | manipulate_data | upload

# Run them. `accounts_pipe` returns a list, so each element starts its own
# `api_pipe` -- three independent chains from one job.
for account in accounts_pipe():
    api_pipe(account)

example.py

import time
import uuid

from buelon.core.step import Result, StepStatus


def accounts(*args) -> list[dict]:
    """The first job in the pipeline. Takes no arguments and returns a list."""
    return [
        {'account_id': 123, 'account': 'mr. business'},
        {'account_id': 456, 'account': 'mrs. business'},
        {'account_id': 789, 'account': 'sr. business'},
    ]


def request_report(account: dict) -> dict:
    """Ask the (imaginary) API for a report. Returns whatever the next job needs."""
    return {**account, 'report_id': f'{uuid.uuid4()}', 'requested_at': time.time()}


def get_status(request: dict) -> Result | dict:
    """Poll until the report is ready.

    `StepStatus.pending` hands the job back to the hub, which re-queues it, so this
    job runs again later instead of blocking a worker for the whole wait.
    """
    if time.time() - request['requested_at'] < 10:
        return Result(status=StepStatus.pending)
    return request


def get_report(request: dict) -> list[dict]:
    """Download the report. Returns a table -- a list of flat dicts."""
    return [
        {**request, 'sales': i * 10.0, 'spend': i * 3.0}
        for i in range(1, 50)
    ]


def upload_to_db(table: list[dict]) -> None:
    """The last job. Returning None is fine."""
    print(f'uploaded {len(table)} rows for {table[0]["account"]}')

Syntax notes

  • Indentation is four spaces. Change it with TAB = ' ' on the first line.
  • !scope, !priority, !timeout, !retries and !max_handbacks set defaults for the whole file when they are at the left margin, and override them for one job when they are indented inside a job definition or attached to an import (...) entry.
  • !retries N gives a failed job N further attempts. They are spaced out, not immediate: the hub holds the job back for 5s, then 10s, then 20s and so on, capped at 5 minutes, so a rate limit or a failover has time to clear. bue status counts the jobs currently waiting as delayed. Tune it with BUELON_RETRY_BACKOFF_BASE / _MAX on the hub.
  • !max_handbacks N caps how many times a job may return StepStatus.pending before the hub gives up and records it as an error (bue run -f fails the job instead). It defaults to 0, meaning unlimited, because a poll loop genuinely does not know how many turns it needs — set it only on a job that should not poll forever. It is a separate budget from !retries, which counts failures.
  • !timeout takes an arithmetic expression in seconds (20 * 60, 60**2 * 5), but it must not contain parentheses inside an import (...) block — the parser counts brackets.
  • A single-job pipe needs a leading |: p = | accounts.
  • A pipe can be wrapped across lines in parentheses.
  • Only two ways to run a pipe: pipe() on its own, or for x in pipe1(): pipe2(x).
  • bue upload runs the script locally and uploads the jobs it produces. A .bue file is a program, not a manifest: bue upload executes it on your machine to build the job graph, then sends the resulting jobs to the hub. Mostly that is just parsing — but a for loop has to know how many jobs to create, so the loop's source pipe genuinely runs locally, in the uploading process and working directory, and outside any !scope. Use bue submit instead to do that build on a worker (see below).
  • # starts a comment.

bue upload vs bue submit

bue upload bue submit
where the script is built your machine a worker, in -s SCOPE
what is sent to the hub the finished jobs one bootstrap job carrying the script
a for loop's source pipe runs locally, no !scope / !timeout / !retries runs as a normal job, with all three
needs the script's imports and referenced files on your machine on the worker
you see build errors immediately, in your terminal in bue errors

submit is the one to reach for when the loop source is expensive, needs credentials or network access your laptop does not have, or belongs on a machine in a particular scope. upload is simpler and tells you about syntax errors on the spot, so it stays the default.

Production Notes

Security. The hub speaks a custom encrypted protocol with no authentication: anything that can reach the port can queue and run arbitrary code. Keep the hub and its workers on a private network, and put anything user-facing (the bue web UI, an upload endpoint) in front of it rather than exposing the hub itself. Set CRYPTO_KEY to a real secret on every process — hub, workers and any machine running bue upload / bue status — or the built-in default key is used and the transport is effectively unencrypted. hub.encryption picks the wire format and must be identical on every process; off disables encryption entirely.

Sizing. One hub, N workers. Each worker runs up to 25 jobs concurrently on an asyncio loop, so the useful number of worker processes is driven by how CPU-bound your jobs are; for the I/O-heavy work Buelon is built for, a handful of processes per machine is plenty. Use scopes to route heavy jobs to the machines that can take them.

Restarts. Workers are disposable — a worker that dies mid-job has its jobs requeued by the hub, and it exits by itself every 20 minutes anyway, so run it under a supervisor. The hub is not disposable: it holds the queue and every job result, and loses up to BUELON_AUTO_SAVE_INTERVAL seconds of progress on an unclean stop.

Memory. The hub keeps every intermediate result until the whole DAG finishes, and a pipeline parked on an error keeps its chain's results for as long as the error sits there. That is deliberate, not a leak: bue reset requeues the errored job by itself, so its parents' results have to still be there when you fix the code two days later and re-run it. It is also what the web UI's job tree shows you when you click into a failure.

The cost is that those results are re-serialized on every autosave. bue status and the web UI report it, so it is a number you can watch rather than a surprise:

$ bue status
done: 0, queued: 18, errors: 4, jobs: 3, delayed: 0, holds: 0, remaining: 21, total: 21, results: 12 (~4.7 MB), staged: 0 in 0 upload(s), handed back: 0 (max 0)

results counts held job results and is deliberately outside total — a result is not a job. The size is estimated from a sample, hence the ~. Two things release it: the pipeline finishing (a DAG that fully succeeds drops all of it), or bue delete, which discards errored pipelines outright. bue reset only releases it if the re-run succeeds. Errors are per pipeline, so one parked chain does not hold another one's results.

staged is the other number outside total: the chunks of a bue upload that has not finished yet. A multi-chunk upload is buffered on the hub and only enters the queue when the whole thing commits, so those jobs are not runnable and must not count towards total — but the hub is holding them in memory, and a non-zero staged that never moves is a stalled uploader. An abandoned buffer is discarded when the connection drops, or reaped after fifteen minutes if the client hangs around without sending anything.

handed back is the third: how many jobs still in play have returned pending at least once, and the highest count among them. These are already counted in jobs and holds — the number exists because a job polling an API that will never be ready looks exactly like a job waiting its turn, and max 4,000 on an otherwise quiet hub is the only thing that gives it away. There is no cap by default; add !max_handbacks N to a job that should give up and land in bue errors instead of polling forever.

Known Defects

  • The postgres: block in settings.yaml is not read by anything. Postgres jobs use the POSTGRES_* environment variables instead.
  • bue demo starts a bucket server that nothing uses and is not a useful demo.
  • Error handling and logging work but are thin.

Future Plans

If this project sees some love, or I just find more free time, I'd like to support more languages like javascript and even compiled languages such as rust, go and c++, allowing teams that write different languages to work on the same program.

Web app for logging, execution and worker management.

Add a scheduler process to allow scheduled pipelines.

Create an official programming/scripting language for parallel processing. This would be separate from the current DML while still being designed to use the Buelon orchestration system.

In Loving Memory

In loving memory of Buelon Rexford Moss.

License

  • MIT License

Metadata

Release files for buelon 1.0.79a1

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for buelon 1.0.79a1
File Size Uploaded
buelon-1.0.79a1.tar.gz 1.5 MB Details

Built distribution (wheel)

Table of built distributions (wheels) for buelon 1.0.79a1
File Interpreter ABI Platform
buelon-1.0.79a1-py3-none-any.whl Python 3 none any Details

Total release size: 3.0 MB

Release files / buelon-1.0.79a1.tar.gz

Download URL buelon-1.0.79a1.tar.gz
Size 1.5 MB
Tags Source
SHA-256 checksum
How to use checksums
ca182aff089bc23343f18647924fb4f6fa19ab92f80e5cfa62b79cf5081c555f
BLAKE2b-256 checksum
How to use checksums
ee75c3ac8cb74e696946f1b1f0f7cfe2fba2fef9be74bea5d82f6a8f0daf5757
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.13.7

Release files / buelon-1.0.79a1-py3-none-any.whl

Download URL buelon-1.0.79a1-py3-none-any.whl
Size 1.5 MB
Tags Python 3
SHA-256 checksum
How to use checksums
d89db6c7debc0de659d1af3113d58935f1bcddff0ba2bf290c5fdc2a096fb9cc
BLAKE2b-256 checksum
How to use checksums
7441834a0d4f2755ea49c23e2d8d9292b4ecd747e031f57dd090f3a52d470aab
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.13.7

Release history Release notifications | RSS feed

This release

1.0.79a1 This release

2 release files

1.0.74

2 release files

1.0.73

2 release files

1.0.68

2 release files

1.0.66

2 release files

1.0.65

2 release files

1.0.64

2 release files

1.0.63

2 release files

1.0.62

1 release file

1.0.61

2 release files

1.0.60

2 release files

1.0.58

2 release files

1.0.40

2 release files

1.0.39

2 release files

1.0.38

2 release files

1.0.37

2 release files

1.0.36

2 release files

1.0.35

2 release files

1.0.34

2 release files

1.0.33

2 release files

1.0.32

2 release files

1.0.31

2 release files

1.0.30

2 release files

1.0.29

2 release files

1.0.28

2 release files

1.0.27

2 release files

1.0.26

2 release files

1.0.25

2 release files

1.0.24

2 release files

1.0.23

2 release files

1.0.22

2 release files

1.0.21

2 release files

1.0.20

2 release files

1.0.19

2 release files

1.0.18

2 release files

1.0.17

2 release files

1.0.16

2 release files

1.0.15

2 release files

1.0.14

2 release files

1.0.13

2 release files

1.0.12

2 release files

1.0.11

2 release files

1.0.10

2 release files

1.0.9

2 release files

1.0.8

2 release files

1.0.7

2 release files

1.0.6

2 release files

1.0.5

2 release files

1.0.4

2 release files

1.0.3

2 release files

1.0.2

2 release files

1.0.1

2 release files

1.0.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page