A unified Python framework for distributed systems and multi-agent applications over Eclipse Zenoh
Project description
Istos
Decorator-first Python services on Eclipse Zenoh — RPC, streaming, duplex channels, work queues, and durable pub/sub.
Istos is for building clean-architecture software and AI agents. A small, explicit decorator surface keeps your codebase maintainable and right-sized — modules that stay easy to analyze, debug, and develop, including with vibe coding, where the machine writes more of the code and readable, well-bounded components are what keep it under control. It maps network operations onto Python decorators so you can wire request/reply, token streams, interactive agent sessions, jobs, and events without standing up a message broker.
Why Istos
Distributed systems are hard enough without spending your design budget on message brokers. Istos lets you reason about services that answer requests and events that fan out — you draw the boxes and arrows; Istos is the arrows. A two-service prototype and a hundred-service fleet are the same mental model, just more decorators.
There is no Kafka, RabbitMQ, NATS, or Redis Streams cluster to provision.
Durability lives in the peers — via Eclipse Zenoh's
advanced pub/sub — not in a central log you have to operate. pip install,
write a decorated function, run().
Key Features
- Decorators first:
@handle,@stream,@channel,@publish,@subscribe. - Streaming RPC:
@streamyields chunks;stream_query/@stream_clientconsume them. - Duplex channels:
@channel+open_channel/@channel_clientfor multi-turn agents (WebSocket or fabric). - HTTP gateway:
Istos(http_port=8080)plushttp=/ws=— JSON, SSE, WebSocket. Co-host inside FastAPI withistos.asgi.lifespan. Optional MCP tools from@handle. - Smart selectors: Query params like
?limit=5land on your Python arguments. - Schema validation: Type hints / Pydantic at the edge.
- Retry policies:
retry=5(or aRetryPolicy) on queries and subscribers. - Brokerless durability:
durable=Truereplay caches;persist="s3://…"/app.replay(...)when the producer itself can die. - Security: TLS/mTLS at the transport;
TokenAuthorizer/JWTAuthorizer/require_roleson handlers. Default is still open — userequire_auth=Truewhen you mean it. - Pluggable storage: in-memory, Redis, or SQLAlchemy (bring your async driver).
The Mental Model
@handle&@query: 1-to-1 RPC@stream: 1-to-1 streaming RPC (SLM/LLM tokens)@channel: full-duplex sessions (agents)@publish&@subscribe: 1-to-many events@on_liveliness: node discovery & health
Installation
This project uses modern Python packaging via uv.
# Standard installation
uv pip install istos
# Or install from source:
git clone https://github.com/0x416d6972/Istos.git
cd istos
uv pip install -e .
# Or with the optional backends (Redis, SQLAlchemy, S3 persistence, JWT, OTel):
uv pip install -e ".[all]"
Quick Start
1. Registering a Handler
Handlers sit on the network and respond to incoming queries. Istos automatically parses query parameters into your function's arguments.
from istos import Istos
istos = Istos()
@istos.handle(prefix="robot/move")
async def move(distance: int, speed: str = "normal"):
"""
Called when a Zenoh Query hits 'robot/move'.
E.g., Querying 'robot/move?distance=10&speed=fast' automatically binds:
distance=10, speed='fast'
"""
return {"status": "success", "distance": int(distance), "speed": speed}
if __name__ == "__main__":
# Blocks and listens for queries
istos.run()
2. Querying the Network
You can easily query handlers registered anywhere on the Zenoh network using kwargs to build Zenoh Selectors.
from contextlib import asynccontextmanager
from istos import Istos
# Using @istos.query makes it a callable network function
@istos.query("robot/move")
async def query_robot(result):
return result
# Queries run over the service's shared Zenoh session, so trigger them once the
# service is running — e.g. from a handler or a lifespan hook:
@asynccontextmanager
async def on_start(app):
reply = await query_robot(distance=15, speed="fast")
print(f"Robot replied: {reply}")
yield
istos = Istos(lifespan=on_start)
if __name__ == "__main__":
istos.run()
There is no longer a per-call "transient session" fallback: the whole service shares one Zenoh session opened by
istos.run()/run_async(). Calling a@query/@publishbefore the service is running raises a clearRuntimeError.
3. Publishing & Subscribing (Event-Driven)
React to real-time events efficiently.
from istos import Istos
istos = Istos()
# --- Subscriber ---
@istos.subscribe("drone/telemetry")
def on_telemetry(data):
# Triggered automatically when data is pushed to "drone/telemetry"
print(f"Received telemetry via network: {data}")
# --- Publisher ---
@istos.publish("drone/telemetry")
def get_telemetry():
# The return value is automatically published to the network!
return {"battery": 85, "altitude": 120}
if __name__ == "__main__":
# Call the wrapped publisher function to publish the result
get_telemetry()
# Or publish arbitrary data independently
# await istos.publish_once("drone/telemetry", {"battery": 80})
istos.run()
Brokerless durability. Add durable=True and a subscriber that joins late
still receives everything — replayed peer-to-peer from the producer's cache, with
no Kafka/NATS broker to run:
@istos.publish("orders/created", durable=True, cache=1000) # producer keeps a replay log
async def created(order): return order
@istos.subscribe("orders/created", durable=True, replay=1000) # replays history + recovers
async def on_created(event): ...
The producer is the log (Zenoh advanced pub/sub); pair it with the idempotency ledger for effectively-once processing. See Brokerless Durable Messaging.
4. Liveliness Tracking (Heartbeats)
Detect when nodes connect or drop off the network without polling.
# Announce that this node is alive on the network
istos.declare_liveliness("robot/camera1")
# Listen to the network for connection state changes
@istos.on_liveliness("robot/**")
def status_changed(key_expr: str, is_alive: bool):
if is_alive:
print(f"Node connected: {key_expr}")
else:
print(f"ALERT: Node crashed/disconnected -> {key_expr}")
5. One-Shot Commands & State Clearing
Use raw async functions when you want to act imperatively rather than relying on events.
# Quickly shoot out a piece of data
await istos.publish_once("fast/data/pulse", {"system": "ok"})
# Clear/erase network states, especially useful if using persistent StoragePlugins
await istos.delete_once("robot/cache/old_logs")
6. High-Performance Shared Memory (Zero-Copy)
When sending large data arrays (like HD video frames) between handlers on the same hardware, enable POSIX shared memory allocations to avoid copies.
@istos.publish("video/feed", use_shm=True)
def send_frame():
return large_data_array
# The framework automatically manages Zenoh ShmProviders natively!
7. Dependency Injection & Pluggability
Set global defaults on startup, or override them per-endpoint for polyglot persistence and sagas:
from istos import Istos
from istos.consistency import InMemoryStoragePlugin, SqlAlchemyStoragePlugin
istos = Istos(
storage=InMemoryStoragePlugin(), # Global default
)
Each decorator can use its own serializer or storage plugin:
from istos.messages.serialization import MsgPackSerializer
# JSON and Memory by default
@istos.handle("robot/move")
async def move(distance: int): ...
# MsgPack and a durable SQL ledger specifically for this endpoint
@istos.handle("sensor/data", serializer=MsgPackSerializer(),
storage=SqlAlchemyStoragePlugin("postgresql+asyncpg://user:pass@db/istos"))
async def sensor(data): ...
Dependency injection with Depends. @handle, @subscribe, @publish, @query, and @on_liveliness can all declare dependencies that Istos resolves per invocation — plain callables, async callables, sub-dependencies, and yield dependencies with setup/teardown. Use it as a default value or inside Annotated:
from typing import Annotated
from istos import Istos, Depends
istos = Istos()
async def get_db():
conn = await open_conn()
try:
yield conn # torn down after the handler replies
finally:
await conn.close()
@istos.handle("orders/create")
async def create(order_id: int, db: Annotated[Conn, Depends(get_db)]):
return await db.insert(order_id)
Dependencies are cached per invocation (use_cache=True by default), sync dependencies are offloaded to a thread so they don't block the event loop, circular dependencies raise DependencyCycleError, and tests can swap any dependency via istos.dependency_overrides[get_db] = fake_db.
Streaming note: on
@subscribe/@publish, dependencies resolve per message — so ayielddependency opens/closes its resource on every message. For an expensive, long-lived resource (DB pool, sensor socket), create it once in thelifespanand inject the shared instance with a cheapDependsthat just returns it; reserve per-messageyielddeps for genuinely per-message setup.
8. Retry Policies
Add automatic retries with exponential backoff to any query or subscriber. Pass a simple integer or a full RetryPolicy for fine-grained control.
from istos.core.retry import RetryPolicy
# Simple — retry up to 5 times with default exponential backoff
@istos.query("weather/forecast", retry=5)
def get_forecast(result):
return result
# Subscriber with retries — if processing crashes, it retries before giving up
@istos.subscribe("sensor/readings", retry=3)
def on_reading(data):
save_to_database(data) # retried automatically on transient failures
# Advanced — full control over backoff timing and failure handling
@istos.query("weather/forecast", retry=RetryPolicy(
max_retries=10,
delay=1.0,
backoff_factor=3.0,
on_failure=lambda e: print(f"Dead letter: {e}")
))
def get_forecast(result):
return result
9. Schema Validation
Istos automatically validates and type-coerces incoming parameters at the network boundary — before your business logic runs. Supports three modes:
from pydantic import BaseModel
# Mode 1: Type hints → auto-coercion
# Zenoh sends distance="15" (string), Istos casts it to int(15)
# Zenoh sends distance="hello" → rejected with a validation error reply
@istos.handle("robot/move")
async def move(distance: int, speed: str = "normal"):
return {"moved": distance, "speed": speed}
# Mode 2: Pydantic BaseModel → full schema validation
class MoveRequest(BaseModel):
distance: int
speed: str = "normal"
@istos.handle("robot/move")
async def move(request: MoveRequest):
# request is a fully validated Pydantic object with defaults applied
return {"moved": request.distance}
# Mode 3: No type hints → passthrough (backward compatible)
@istos.handle("robot/echo")
async def echo(message):
return {"echo": message}
10. Security: Transport, Authentication & Authorization
[!WARNING] Istos is unauthenticated by default. With no configuration, a session runs in Zenoh peer mode with multicast scouting and no TLS — any peer on the local network can discover your node, invoke every
@handle, and read every published value. Istos raises anIstosSecurityWarning(viawarnings.warn, like urllib3'sInsecureRequestWarning) whenever it opens such a session. You can escalate it to a hard failure in CI:import warnings from istos import IstosSecurityWarning warnings.simplefilter("error", IstosSecurityWarning) # insecure config -> exceptionBefore deploying, do both of the following:
- Secure the transport — authenticate and encrypt the fabric (below).
- Authorize handlers — gate who may invoke them (further below).
Transport Security & Authentication
Secure the system without hard-coding secrets. IstosZenohConfig loads properties from your environment (.env file or environment variables) using pydantic-settings.
# .env file
ISTOS_ZENOH_MODE=client
# Endpoints accept a JSON array or a comma-separated list:
ISTOS_ZENOH_CONNECT_ENDPOINTS=["tls/zenoh-router.local:7447"]
# ISTOS_ZENOH_CONNECT_ENDPOINTS=tls/router-a:7447,tls/router-b:7447
# Basic Auth
ISTOS_ZENOH_USERNAME=robot_1
ISTOS_ZENOH_PASSWORD=super_secret
# TLS / mTLS
ISTOS_ZENOH_ROOT_CA_CERTIFICATE=/path/to/ca.pem
# ISTOS_ZENOH_LISTEN_CERTIFICATE=/path/to/cert.pem
# ISTOS_ZENOH_LISTEN_PRIVATE_KEY=/path/to/key.pem
# ISTOS_ZENOH_ENABLE_MTLS=true
# Lock down discovery: disable UDP multicast scouting and rely on explicit endpoints
# ISTOS_ZENOH_MULTICAST_SCOUTING=false
IstosZenohConfig is a Pydantic BaseSettings model, so it validates at construction (via field/model validators) and fails fast on structurally broken security config (malformed endpoints, a TLS cert without its key, mTLS without a CA, an invalid mode), warns when credentials are configured without TLS, and rejects unknown ISTOS_ZENOH_* variables so a typo can't silently disable auth. Secrets (password, listen_private_key) are SecretStr, so they don't leak into logs or reprs. For knobs the builder doesn't model, deep-merge raw Zenoh config via additional_config={...}. Session managers also accept open_retries/open_retry_delay_s to wait out a router that isn't up yet at startup.
Then pass it straight to Istos — it builds the session for you:
from istos import Istos
from istos.communication.sessions import IstosZenohConfig
# Auto-populates from .env!
config = IstosZenohConfig()
istos = Istos(config=config)
config= accepts an IstosZenohConfig (built via .build() at construction) or a raw zenoh.Config. It's mutually exclusive with session_manager= — pass one or the other. If you need a sync session or full control, wire it yourself instead:
from istos.communication.sessions import AsyncZenohSession
istos = Istos(session_manager=AsyncZenohSession(config.build()))
Advanced Security: Vault & Secret Managers (Programmatic Raw PEM)
For zero-trust environments, Zenoh accepts raw multiline PEM strings natively, so you don't need to write files to disk. You can bypass .env and pull secrets into IstosZenohConfig at startup:
# 1. Fetch from your secrets manager (HashiCorp Vault, AWS, LDAP, etc.)
secrets = vault.get_secret("istos/prod")
# 2. Inject raw strings directly into the config builder
config = IstosZenohConfig(
mode="client",
connect_endpoints=["tls/zenoh-router.local:7447"],
username=secrets["zenoh_user"],
password=secrets["zenoh_pass"],
root_ca_certificate=secrets["raw_ca_pem_string"], # raw multiline PEM
)
istos = Istos(session_manager=AsyncZenohSession(config.build()))
Authorization
Securing the transport controls who joins the fabric. Authorization controls what a joined peer may invoke. Istos gives every handler an authorization hook. An authorizer receives an AuthContext (the key expression, parameters, and any attachment/token the caller sent) and returns True to allow the request, or False/raises UnauthorizedError to deny it. Sync and async authorizers are both supported. Denied requests never reach your handler and are answered with an unauthorized error.
from istos import Istos, TokenAuthorizer, AuthContext, UnauthorizedError
# App-wide: every handler (including built-in .istos/* endpoints) requires a token
istos = Istos(authorizer=TokenAuthorizer("super-secret-token"))
# Per-handler override — a custom rule for one endpoint
def admins_only(ctx: AuthContext) -> bool:
return ctx.token in {"alice-key", "bob-key"}
@istos.handle("fleet/shutdown", authorizer=admins_only)
async def shutdown():
return {"stopping": True}
Callers attach their token via attachment=:
await istos.query_once("fleet/shutdown", attachment="alice-key")
Built-in endpoints (
.istos/health,.istos/ready,.istos/metrics) andserve_docs()inherit the app-wide authorizer. If you leave the app unauthenticated, Istos warns that these endpoints — including the AsyncAPI document that describes your entire API surface — are network-reachable by any peer. Set anauthorizer(or pass one toserve_docs(authorizer=...)) to protect them.
A note on serialization
Istos ships JsonSerializer (default), RawSerializer (bytes/str passthrough for pre-encoded or binary payloads), MsgPackSerializer, PydanticSerializer, ProtobufSerializer, YamlSerializer, and a Base64Serializer wrapper. JsonSerializer tolerates common types like datetime, Decimal, and UUID (stringified), and MsgPackSerializer pins string decoding for cross-version consistency.
It does not ship a pickle-based serializer: pickle.loads executes arbitrary code embedded in its input, which on a fabric where any peer can publish to a key is remote code execution by design.
Testing
Istos includes IstosTestClient for in-process testing without a Zenoh network:
from istos import Istos
from istos.testing import IstosTestClient
istos = Istos()
@istos.handle("robot/move")
async def move(distance: int):
return {"moved": distance}
client = IstosTestClient(istos)
result = await client.query("robot/move", distance=10)
Run the full test suite:
uv pip install -e ".[dev]"
pytest tests/ -m "not integration" # unit tests
pytest tests/ # includes network integration tests
Production Features
- Logging — text or JSON under
istos.*. We don't reconfigure your root logger unless you ask (configure_logging=True/configure_logging()). - Health —
.istos/health/.istos/ready(and/livez//readyzwithhttp_port) - Metrics —
.istos/metrics(and/metrics) - Capabilities —
.istos/capabilitieslists handlers/streams with schemas - HTTP / SSE —
http_port+http=on handle/stream - Tracing — optional OTel (
istos[otel]) - Middleware — correlation IDs, logging, your own
- Exceptions —
@exception_handler - Shutdown — SIGINT/SIGTERM
- Storage — Redis, SQLAlchemy, S3 persist, JWT extras
CLI
istos new my-service # Scaffold a new project
istos docs # Serve documentation locally
istos version # Print installed version
Contributing
Contributions and pull requests are welcome! Ensure tests pass and type hints are satisfied.
- Fork the repository
- Create your feature branch (
git checkout -b feature/amazing-feature) - Commit your changes
- Push to the branch (
git push origin feature/amazing-feature) - Open a Pull Request
License: Apache-2.0
Python: 3.10, 3.11, 3.12, 3.13, and 3.14
Project 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 istos-0.1.0.tar.gz.
File metadata
- Download URL: istos-0.1.0.tar.gz
- Upload date:
- Size: 110.7 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via: uv/0.11.28 {"installer":{"name":"uv","version":"0.11.28","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
64686b90bb5cd02c4da7c051ae31211930d4c1b0f82c25ffee8c1cb5ac945701
|
|
| MD5 |
c3cad33c899c738457a2334c420bc680
|
|
| BLAKE2b-256 |
a6bd94e1583aac46cd626b981818bb138ea0c6ac984b1816e0ecb6f503debd28
|
File details
Details for the file istos-0.1.0-py3-none-any.whl.
File metadata
- Download URL: istos-0.1.0-py3-none-any.whl
- Upload date:
- Size: 143.8 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via: uv/0.11.28 {"installer":{"name":"uv","version":"0.11.28","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
e3f509b7c88a79b373b5f1314ff3ffb25e1072d4fe817b707748e46f547ad0f4
|
|
| MD5 |
00352032fa151810cadf5483ea4dead7
|
|
| BLAKE2b-256 |
ddb4776bb39124189ef642ccef173e1caff76022f74dd64fdc875a8ba6ecba8e
|