Operonx
Operonx is a workflow engine where ops can yield — so the same async DAG handles batch jobs (Airflow-style) and event-driven streaming pipelines (pipecat-style callbot / voice / STT → LLM → TTS).
The Rust execution backend now lives in its own repo: batman1m2001-cyber/operonx-rs (crates.io). It shares the shared JSON spec fixtures with this repo but ships independently.
Why Operonx
- Yield-based streaming. Generator ops emit per-item; downstream dispatches per-frame, not per-batch. The
for_loop/map_op/ VAD → STT → LLM → TTS shapes work without bolt-on map/reduce ops. - Operator reference syntax.
op["key"],PARENT["key"],op["src"] >> PARENT["dst"],outputs={"*": PARENT}— explicit and local. Noxcom_pullper node, no JSON serialisation per hop. - Multi-provider LLM / embedding / rerank. OpenAI, Azure, Gemini, Anthropic, vLLM, TEI, HuggingFace, ONNX, Pinecone — swap with one line in
resources.yaml. Built-in weighted load balancing + fallback chains. - Tracing built-in. Langfuse, OpenTelemetry, and a local file consumer. All async-flushed; never blocks the run.
- Lean tier-1.
pip install operonxis justpydantic / pyyaml / rich / orjson. Provider SDKs are extras.
Quick Start
pip install operonx
import asyncio
from operonx.core import Operon, GraphOp, op, START, END, PARENT
@op
def greet(who: str):
return {"message": f"Hello, {who}!"}
async def main():
with GraphOp(name="hello") as graph:
step = greet(who=PARENT["who"])
START >> step >> END
result = await Operon(graph).run(inputs={"who": "World"})
print(result["message"]) # Hello, World!
asyncio.run(main())
Streaming with yield
The differentiator. A generator op yields per item; downstream ops dispatch on each frame. The same engine that runs a batch DAG runs a callbot pipeline.
from operonx.core import Operon, GraphOp, op, START, END, PARENT
@op
def chunk_text(text: str, chunk_size: int):
for i, words in enumerate(words_in(text, chunk_size)):
yield {"chunk": " ".join(words), "index": i}
@op
def analyze(chunk: str, index: int):
return {"result": f"[{index}] {len(chunk.split())} words"}
with GraphOp(name="pipeline") as g:
src = chunk_text(text=PARENT["text"], chunk_size=PARENT["chunk_size"])
step = analyze(chunk=src["chunk"], index=src["index"])
START >> src >> step >> END
Each yield triggers a dispatch on a fresh (parent_ctx, "yield_N") sub-context. Empty yield = zero downstream dispatches (matches Python's skipped yield). N-to-M flows (one VAD chunk → multiple speech segments) work because each yield is independent.
See examples/python/ex14 for the streaming + tracing demo, examples/python/ex15 for the callbot pipeline (audio → VAD → STT → intent → handler → TTS).
LLMs in one line
pip install "operonx[standard]"
import asyncio
import operonx
from operonx.core import Operon, GraphOp, START, END, PARENT
from operonx.providers import LLMOp
async def main():
operonx.bootstrap() # loads ./.env + ./resources.yaml
with GraphOp(name="qa") as graph:
c = LLMOp(
name="llm",
resource="gpt-4o-mini",
inputs={
"prompt": {"system": "You are a helpful assistant.", "user": "{question}"},
"*": PARENT,
},
outputs={"*": PARENT},
)
START >> c >> END
result = await Operon(graph).run(inputs={"question": "What is Python?"})
print(result["content"])
asyncio.run(main())
LLMOp.prompt accepts a string, {"system": ..., "user": ...} dict, or a full messages list — every non-reserved kwarg becomes a {var} substitution.
Multi-model load balancing + fallback
from operonx.providers import LLMOp
llm = LLMOp.of(
resource=["gpt-4o", "gpt-4o-mini"],
ratios=[0.7, 0.3], # 70 / 30 split
fallback=["claude-haiku"], # tried in order on failure
messages=PARENT["messages"],
)
Branching
from operonx.core import START, END, GraphOp, PARENT
from operonx.core.ops.flow.branch_op import if_
router = (if_(PARENT["score"] >= 90, "excellent")
.if_(PARENT["score"] >= 70, "good")
.else_("fail"))
START >> router >> excellent >> merge >> END
router >> good >> merge
router >> fail >> merge
if_() evaluates conditions in order; the first match routes through a soft edge (>>~ semantically — branch outputs use soft edges so non-matching branches don't block downstream).
Loops
Write a back-edge inside @graph — the build-time cycle-rewrite pass
turns it into a hidden _GraphLoop so the scheduler still sees a DAG:
from operonx.core import graph, START, END, PARENT
from operonx.core.ops.flow.branch_op import if_
@graph
def counter():
PARENT.declare(count=0)
inc = increment(counter=PARENT["count"])
inc["counter"] >> PARENT["count"]
START >> inc >> if_(PARENT["count"] >= 5, END).else_(inc)
g = counter()
The branch's else_ target is the back-edge; each iteration commits
its outputs to the shared count cell and the branch decides whether
to loop again or exit. See docs/guide/03-loops-and-branches.md.
Installation
Single Python package, optional extras for each integration:
pip install operonx # Tier 1 — engine only, ~10 MB
pip install "operonx[openai]" # OpenAI / Azure
pip install "operonx[anthropic]" # Anthropic via httpx
pip install "operonx[gemini]" # Vertex AI
pip install "operonx[onnx]" # Local ONNX inference
pip install "operonx[langfuse]" # Langfuse tracing
pip install "operonx[otel]" # OpenTelemetry tracing
pip install "operonx[standard]" # Recommended — providers + Langfuse + OTEL
pip install "operonx[all]" # Everything except torch / HuggingFace
| Extra | Contents |
|---|---|
openai |
OpenAI SDK (also covers Azure) |
anthropic |
httpx + OpenAI message types |
gemini |
google-cloud-aiplatform + AsyncOpenAI client |
bedrock |
boto3 + OpenAI message types |
onnx |
onnxruntime + tokenizers + numpy |
huggingface |
transformers + torch (~2.5 GB; opt in) |
langfuse |
Langfuse SDK |
otel |
OpenTelemetry API + SDK + OTLP exporters |
standard |
OpenAI + Langfuse + OTEL (production bundle) |
all |
Every provider + tracer except huggingface |
dev |
pytest, ruff, pre-commit |
Tracing
import operonx
from operonx.core import Operon
operonx.bootstrap() # registers consumer configs from resources.yaml
engine = Operon(graph, trace=["trace_langfuse:default"])
Consumers are configured in resources.yaml (trace_local:, trace_langfuse:) and referenced by key. See docs/api/telemetry.md for the full V3 tracing API.
Documentation
| Need | Go to |
|---|---|
| Runnable examples (Python) | examples/python/ |
| Architecture | docs/architecture/ |
| User guide | docs/guide/ |
| API reference | https://batman1m2001-cyber.github.io/Operonx/ |
| Rust runtime | operonx-rs |
Contributing
git clone https://github.com/batman1m2001-cyber/Operonx.git
cd Operonx
uv sync --all-extras
pre-commit install
uv run pytest tests/ -m "not integration"
See CONTRIBUTING.md for the full contributor guide.
License
Apache 2.0
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 operonx-1.0.0.tar.gz.
File metadata
- Download URL: operonx-1.0.0.tar.gz
- Upload date:
- Size: 219.9 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
51232aa5a8f86af3953087dbd12779d8d3501334846e6dbd961827bcd5927723
|
|
| MD5 |
dd2a420f540c9c50d8f7c58d05b8e695
|
|
| BLAKE2b-256 |
bb3885efcecb66a5b572a80c67e9e2854804a5232be22c3a742021bb7ebc0f25
|
File details
Details for the file operonx-1.0.0-py3-none-any.whl.
File metadata
- Download URL: operonx-1.0.0-py3-none-any.whl
- Upload date:
- Size: 288.4 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
f913720e15cc26d7b751b90682a55b5c48e68647715c846f8ee94ea7054597e6
|
|
| MD5 |
d00a53136ff0078b33e42fce82547885
|
|
| BLAKE2b-256 |
80d3e454a31a7d3ab8a1be4568a81f94e8547096404cb586b294e64288e0db12
|