Cellaflow Python SDK
The official Python SDK for the CellaFlow Engine — providing durable execution, deterministic replay recovery, and swarm-safe concurrency primitives for AI workflows and multi-agent systems.
Features
- 🔄 Durable Execution & Transparent Replay: Workflows survive process restarts and infrastructure crashes without re-executing completed steps.
- ⚡ Zero-Friction Decorators: Annotate standard Python functions with
@workflow,@step, and@tool(supporting bothdefandasync def). - 🧠 LangGraph Drop-in Checkpointer: Native
CellaflowSavercheckpointer for LangGraph workflows with immutable MessagePack state snapshots, time-travel, and crash recovery. - 🛡️ Deterministic Idempotency: Automatic input hashing via RFC 8785 Canonical JSON and SHA-256 guarantees cross-language determinism across multi-agent swarms.
- 🔒 Background Lease Management: Non-blocking heartbeat management keeps engine locks alive and eliminates split-brain execution using fencing tokens.
- 📦 Secure Serialization: Strictly uses MessagePack for all state payloads to optimize throughput over gRPC and eliminate Remote Code Execution (RCE) deserialization vectors.
Installation
Core SDK
For standard durable execution workflows, decorators, and gRPC client:
pip install cellaflow
With LangGraph Support
To use CellaflowSaver as a drop-in LangGraph checkpointer:
pip install "cellaflow[langgraph]"
Note on Optional Dependencies:
pip install cellaflowinstalls a lightweight package without pulling in LangGraph, Pydantic, or LangChain.- If your environment already contains
langgraph >= 0.2.0andlanggraph-checkpoint >= 2.0.0,from cellaflow import CellaflowSaverwill work immediately out of the box.- Installing
cellaflow[langgraph]ensures that compatiblelanggraphandlanggraph-checkpointversions are installed or upgraded automatically.
Quickstart
1. Synchronous Workflow
from cellaflow import workflow, step, tool, IdempotencyScope
# 1. Define tools with automatic caching and idempotency
@tool(scope=IdempotencyScope.SCOPE_SESSION_WIDE)
def search_web(query: str) -> dict:
"""Simulates an expensive API or web search tool."""
print(f"Executing web search for: {query}")
return {"query": query, "results": ["CellaFlow Overview", "Durable Execution Docs"]}
# 2. Define intermediate steps
@step
def summarize_results(data: dict) -> str:
return f"Processed {len(data['results'])} items for '{data['query']}'"
# 3. Define the orchestrating workflow
# By default connects to localhost:50051 (or pass custom target="engine:50051", secure=True)
@workflow(version="1.0.0")
def research_workflow(topic: str) -> str:
search_data = search_web(topic)
summary = summarize_results(search_data)
return summary
if __name__ == "__main__":
result = research_workflow("autonomous agent architectures")
print("Result:", result)
2. Asynchronous Workflow & Multi-Agent Swarms
The SDK natively supports async def coroutines, asynchronous I/O, and agent isolation:
import asyncio
from cellaflow import workflow, step, tool, IdempotencyScope
@tool(agent_id="researcher_agent", scope=IdempotencyScope.SCOPE_AGENT_PRIVATE)
async def fetch_market_data(ticker: str) -> dict:
print(f"Fetching live ticker data: {ticker}")
await asyncio.sleep(0.5) # Simulate non-blocking async network I/O
return {"ticker": ticker, "price": 142.50}
@step(agent_id="analyst_agent")
async def analyze_trends(market_data: dict) -> str:
return f"Signal for {market_data['ticker']}: BUY at ${market_data['price']}"
@workflow(version="1.0.0")
async def swarm_analysis_workflow(ticker: str) -> str:
data = await fetch_market_data(ticker)
analysis = await analyze_trends(data)
return analysis
if __name__ == "__main__":
result = asyncio.run(swarm_analysis_workflow("CELL"))
print("Analysis Result:", result)
3. Transparent Replay & Session Recovery
To recover an interrupted execution after a process crash or restart, supply the _session_id keyword argument:
# Resumes the workflow from the exact step where it stopped.
# Completed steps are loaded from the engine's event graph and returned with 0ms execution time.
recovered_result = research_workflow(
"autonomous agent architectures",
_session_id="3f7491d9-e932-4467-bcda-370fb5c1a7e4"
)
4. LangGraph Checkpointer (CellaflowSaver)
Use CellaflowSaver as a drop-in checkpointer for any LangGraph StateGraph:
from typing import TypedDict, List
from langgraph.graph import StateGraph, START, END
from cellaflow import CellaflowSaver
# 1. Initialize checkpointer connected to CellaFlow Engine
checkpointer = CellaflowSaver(target="localhost:50051", secure=False)
# 2. Define LangGraph state schema
class WorkflowState(TypedDict):
query: str
summary: str
def plan_node(state: WorkflowState) -> dict:
return {"summary": f"Plan generated for {state['query']}"}
# 3. Build graph
builder = StateGraph(WorkflowState)
builder.add_node("plan", plan_node)
builder.add_edge(START, "plan")
builder.add_edge("plan", END)
# 4. Compile with CellaFlow persistence
app = builder.compile(checkpointer=checkpointer)
# 5. Execute with session tracking
config = {"configurable": {"thread_id": "session-42"}}
result = app.invoke({"query": "AI Agents"}, config)
print("Graph output:", result)
# 6. Time travel / inspect state history
for checkpoint_tuple in checkpointer.list(config):
print("Checkpoint ID:", checkpoint_tuple.checkpoint["id"])
print("State snapshot:", checkpoint_tuple.checkpoint["channel_values"])
Architecture & How It Works
1. Transparent Replay Recovery
If a worker crashes mid-workflow, restarting the workflow with the same session automatically recovers state from the CellaFlow engine's durable event log:
flowchart LR
classDef app fill:#e1f5fe,stroke:#0288d1,stroke-width:2px,color:#01579b
classDef engine fill:#ede7f6,stroke:#512da8,stroke-width:2px,color:#311b92
classDef crash fill:#ffebee,stroke:#c62828,stroke-width:2px,color:#b71c1c
subgraph Client[" 🐍 Python Application "]
W["@workflow Orchestrator"]:::app
S1["@step 1: Query API"]:::app
S2["@step 2: Process Data"]:::app
Crash["💥 Process Crashes Mid-Run"]:::crash
end
subgraph Storage[" ⚙️ CellaFlow Engine "]
Log[("💾 RocksDB Event Graph")]:::engine
Replay["🔄 On-Demand Replay Engine"]:::engine
end
W --> S1
S1 -->|"1. Commit Result"| Log
S1 --> S2
S2 -.-> Crash
Crash ==>|"Restart Workflow"| Replay
Replay -->|"2. Instant Cache Replay (0ms)"| S1
Replay -->|"3. Resume Live Execution"| S2
2. Idempotency & Lease Lifecycle
For external tool calls and multi-agent coordination, @tool prevents duplicate tool calls, single-flights in-progress work, and maintains active lock heartbeats:
flowchart TD
classDef check fill:#e0f7fa,stroke:#00838f,stroke-width:2px,color:#004d40
classDef hit fill:#e8f5e9,stroke:#2e7d32,stroke-width:2px,color:#1b5e20
classDef lease fill:#fff3e0,stroke:#ef6c00,stroke-width:2px,color:#e65100
classDef commit fill:#ede7f6,stroke:#512da8,stroke-width:2px,color:#311b92
In["Function Call: @tool(*args, **kwargs)"] --> Hash["🔑 RFC 8785 Canonical JSON + SHA-256 Hashing"]
Hash --> Query{"Check Idempotency Cache"}:::check
Query -->|"⚡ Cache HIT"| Cached["Return Cached StepResult<br/>(Skip Execution)"]:::hit
Query -->|"🔒 Cache ACQUIRED"| Lease["Acquire Fencing Token & Lease"]:::lease
subgraph Execution[" ⚙️ Active Tool Execution "]
Lease --> HB["🔄 Background LeaseHeartbeat<br/>(Periodic RenewLease)"]:::lease
Lease --> Run["🛠️ Run User Function Body"]:::lease
end
Run --> Done["Stop Heartbeat & CommitStep(fencing_token)"]:::commit
Done --> DB[("💾 Persist Result to RocksDB")]:::commit
Core Concepts
@workflow(version="1.0.0", target="localhost:50051", secure=False)
Marks the entry point for a durable workflow execution:
- Automatic Client Management: Transparently initializes the underlying gRPC client and session state.
- Context Isolation: Uses Python's
contextvarsto manage session context safely across async event loops and threads. - Replay State Loader: Loads completed step history from the engine whenever
_session_idis supplied.
@step and @tool
Decorators for atomic units of execution within a workflow:
@step(
idempotency_key: Optional[str] = None,
agent_id: str = "default",
tool_name: Optional[str] = None,
scope: IdempotencyScope = IdempotencyScope.SCOPE_SESSION_WIDE,
)
- Idempotency Key Derivation: Automatically builds a deterministic key using:
[Session_ID]:[Workflow_Version]:[Step_Sequence]:[Agent_ID]:[Tool_Name]:[RFC8785_SHA256_Hash] - Replay Interception: If a step was already completed in the session history, immediately returns the cached output without re-running the function body.
- Lease Heartbeating: Automatically runs background heartbeats (
RenewLease) via daemon threads (sync) or asyncio tasks (async) to keep engine locks refreshed. - Lock Release on Failure: If an unhandled exception occurs, automatically releases the lease with
reason="TOOL_ERROR".
IdempotencyScope
Controls how cached step and tool results are shared across multi-agent sessions:
IdempotencyScope.SCOPE_SESSION_WIDE(Default): Cached results are shared across all agents in the session.IdempotencyScope.SCOPE_AGENT_PRIVATE: Isolates cache hits to the executingagent_id.IdempotencyScope.SCOPE_STEP_LOCAL: Strictly isolates cache hits to the current superstep sequence andagent_id.
Low-Level CellaflowClient
For advanced use cases requiring manual session management, graph inspection, or custom scheduling, use CellaflowClient directly:
from cellaflow import CellaflowClient
from cellaflow.v1.common_pb2 import STEP_STATUS_SUCCESS
# Initialize client
client = CellaflowClient(target="localhost:50051", secure=False)
# 1. Start or resume a session
session_resp = client.start_session(
workflow_id="manual_workflow",
version="1.0.0"
)
session_id = session_resp.session_id
# 2. Inspect session history graph
steps, next_cursor = client.get_graph(session_id=session_id)
# 3. Commit a step result
commit_resp = client.commit_step(
session_id=session_id,
sequence=1,
name="custom_task",
status=STEP_STATUS_SUCCESS,
output_payload={"result": "data"}
)
# 4. Clean up channel
client.close()
Development Setup
To get up and running with the Python SDK for local development:
-
Create and activate a virtual environment:
cd python python3 -m venv venv source venv/bin/activate
-
Install the SDK in editable mode with development dependencies:
pip install -e ".[dev]"
Generating Protobufs and Typing Stubs (.pyi)
We use grpcio-tools and mypy-protobuf to generate Python stubs from .proto definitions:
./scripts/generate_protos.sh
This compiles protos from ../proto/cellaflow/v1/*.proto and outputs *_pb2.py, *_pb2_grpc.py, and *.pyi typing stubs into src/cellaflow/v1/.
Running Tests and Linters
# Run unit tests
pytest
# Formatting & Linting
black src tests
flake8 src tests
mypy src tests
License
Apache 2.0 or MIT (see repository root for details).
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 cellaflow-0.2.0.tar.gz.
File metadata
- Download URL: cellaflow-0.2.0.tar.gz
- Upload date:
- Size: 43.7 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via:
twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
3659d54c02ba1c9d3a6ab0be16a318449bd87fda875691acb12532fe47905ff9
|
|
| MD5 |
24843dd24b9a2cf8e3a30fa3a346ecc3
|
|
| BLAKE2b-256 |
0a5008b565791746d306f96d32e05cad3bac07bf2b08a7c61206693062c5304d
|
Provenance
The following attestation bundles were made for cellaflow-0.2.0.tar.gz:
Publisher:
publish-python-sdk.yml on cellaflow/cellaflow-sdks
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
cellaflow-0.2.0.tar.gz -
Subject digest:
3659d54c02ba1c9d3a6ab0be16a318449bd87fda875691acb12532fe47905ff9 - Sigstore transparency entry: 2568950947
- Sigstore integration time:
-
Permalink:
cellaflow/cellaflow-sdks@6e71605f732c0096dd05908fcfae4ead7099c071 -
Branch / Tag:
refs/heads/main - Owner: https://github.com/cellaflow
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
publish-python-sdk.yml@6e71605f732c0096dd05908fcfae4ead7099c071 -
Trigger Event:
workflow_dispatch
-
Statement type:
File details
Details for the file cellaflow-0.2.0-py3-none-any.whl.
File metadata
- Download URL: cellaflow-0.2.0-py3-none-any.whl
- Upload date:
- Size: 39.9 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via:
twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
5a93078b576bbd6408e19fd556e303886dd678a2fefb6a7163b0473b5ae13a35
|
|
| MD5 |
9ad56ca7ea91a1f488e3dca89bd00e09
|
|
| BLAKE2b-256 |
00960f26d23de137f7bfd5376847de6dc401df2ad535c12d164ebf3fec7c5be5
|
Provenance
The following attestation bundles were made for cellaflow-0.2.0-py3-none-any.whl:
Publisher:
publish-python-sdk.yml on cellaflow/cellaflow-sdks
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
cellaflow-0.2.0-py3-none-any.whl -
Subject digest:
5a93078b576bbd6408e19fd556e303886dd678a2fefb6a7163b0473b5ae13a35 - Sigstore transparency entry: 2568950952
- Sigstore integration time:
-
Permalink:
cellaflow/cellaflow-sdks@6e71605f732c0096dd05908fcfae4ead7099c071 -
Branch / Tag:
refs/heads/main - Owner: https://github.com/cellaflow
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
publish-python-sdk.yml@6e71605f732c0096dd05908fcfae4ead7099c071 -
Trigger Event:
workflow_dispatch
-
Statement type: