High-performance communication middleware for scoped NATS messaging.
Project description
KinopioHub
Cloud-native communication middleware for working with remote variables and functions as scoped local objects.
Features
- Automatic value tracking and local caching
- Hierarchical scope system with dynamic attribute access
- Publish/subscribe and request/reply in one API surface
- Automatic reconnection handling with retry backoff
- Native support for TCP, TLS, WebSocket, and secure WebSocket NATS endpoints
- Pluggable codecs for custom serialization strategies
Installation
pip install kinopio-hub
The default install includes the runtime dependencies used by kinopio_hub.leaf.
The local nats-server binary is still resolved lazily at runtime and is never downloaded during
pip install. The package also installs the kinopio-hub console script for local leaf runtime
operations.
Quick Start
import asyncio
from kinopio_hub import KinopioHub
async def main() -> None:
hub = KinopioHub(
servers=["nats://demo.nats.io:4222"],
debug=True,
)
await hub.wait_connected()
chat_messages = hub.chat.messages
await chat_messages.publish({"user": "alice", "message": "hello"})
async def on_message(data, _message) -> None:
print("Received:", data)
await chat_messages.subscribe(on_message)
asyncio.run(main())
Core Concepts
KinopioHub
KinopioHub manages the NATS connection lifecycle and gives access to scopes.
hub = KinopioHub(
servers=["nats://demo.nats.io:4222"],
debug=True,
no_echo=False,
reconnect_timeout=5.0,
)
Scopes
Scopes group related variables under a shared subject prefix.
users = hub.get_scope("users")
online = users.get_variable("online")
count = users.get_variable("count")
same_online = hub.users.online
Variables
Variables are subject-backed objects that support publishing, subscribing, request/reply, and service handlers.
await hub.system.health.publish({"cpu": 42.1})
async def on_update(data, _message) -> None:
print("Updated value:", data)
await hub.system.health.subscribe(on_update)
response = await hub.math.calculator.request({"operation": "add", "a": 3, "b": 4})
await hub.math.calculator.serve(lambda request, _message: {"result": request["a"] + request["b"]})
Every variable exposes value, which contains the latest known value seen locally.
Configuration
| Option | Type | Default | Description |
|---|---|---|---|
servers |
Sequence[str] |
["nats://demo.nats.io:4222"] |
NATS server URLs |
debug |
bool |
False |
Enable debug logging |
no_echo |
bool |
False |
Do not receive messages published by the same client |
server_selection_mode |
"ordered" | "random" | "latency" |
None |
Explicit multi-server selection strategy |
no_randomize |
bool |
True |
Keep the server list order unchanged |
max_reconnect_attempts |
int |
-1 |
Maximum reconnect attempts for the underlying client |
wait_on_first_connect |
bool |
True |
Start connecting as soon as an event loop is available |
reconnect_timeout |
float |
5.0 |
Initial connection timeout in seconds |
reconnect_time_wait |
float |
0.5 |
Delay between reconnect attempts in seconds |
ping_interval |
int |
3 |
Ping interval for the underlying NATS client |
max_ping_out |
int |
3 |
Maximum unanswered pings before the client reconnects |
timeout |
float |
3.0 |
Default request timeout in seconds |
health_report |
float |
5.0 |
Debug health report interval in seconds |
auto_retry |
bool |
True |
Retry initial connection failures automatically |
retry_delay |
float |
1.0 |
Initial retry delay for the connection loop |
retry_backoff_factor |
float |
1.5 |
Exponential backoff factor for retries |
max_retry_delay |
float |
30.0 |
Upper bound for retry backoff |
codec |
KinopioCodec |
None |
Custom codec with encode() and decode() |
json_default |
Callable |
None |
Fallback hook for json.dumps() |
json_object_hook |
Callable |
None |
Fallback hook for json.loads() |
tls |
ssl.SSLContext | bool |
None |
Custom TLS context, or True to use the default system trust context |
tls_hostname |
str |
None |
Explicit TLS hostname override |
tls_handshake_first |
bool |
False |
Use a TLS-first handshake for servers configured with handshake_first: true |
ws_connection_headers |
Mapping[str, Sequence[str]] |
None |
Extra headers for WebSocket transport |
name |
str |
None |
NATS client name |
no_randomize is kept for compatibility. If server_selection_mode is not set, no_randomize=True
maps to "ordered" and no_randomize=False maps to "random". If
server_selection_mode is set explicitly, it takes precedence. "latency" mode opens a short-lived
connection to each candidate server, measures a flush() round-trip, and then orders candidates by
healthy first -> lower RTT first -> original input order. If every probe fails, the original
server order is preserved and the normal connection error flow continues. After connecting in
"latency" mode with multiple candidates, the client re-probes every 10 minutes and hot-switches
only when another healthy server is at least 30ms faster. During the brief double-subscription
window of a hot switch, a small number of duplicate callback deliveries is possible, but logical
subscriptions, services, and value tracking continue to be maintained.
If your NATS servers are configured with TLS-first handshakes (tls.handshake_first: true on the
server side), also set tls_handshake_first=True. Those endpoints do not send the initial INFO
line in clear text before the TLS upgrade. For public CA certificates, tls=True uses Python's
default system trust context. For self-signed or private CA certificates, pass a custom
ssl.SSLContext.
When pairing this Python package with KinopioHub.JS, keep the transport roles explicit: JS clients use WebSocket or secure WebSocket endpoints, while Python clients use native NATS TCP or TCP/TLS endpoints. Cross-implementation compatibility comes from matching subjects and serialized payloads, not from using the same transport.
Connection Management
await hub.wait_connected()
await hub.reconnect()
await hub.aclose()
To track connection state:
from kinopio_hub import ConnectionState
def handle_state(state: ConnectionState) -> None:
print("state:", state.value)
stop = hub.on_state_change(handle_state)
stop()
Local Leaf Runtime
The separate kinopio_hub.leaf module starts a local leaf runtime without changing the main
KinopioHub export surface.
import ssl
from kinopio_hub import KinopioHub
from kinopio_hub.leaf import LeafNodeOptions, start_leaf_node
with start_leaf_node(
LeafNodeOptions(
backbone_servers=("nats://127.0.0.1:7422",),
)
) as leaf:
print(leaf.client_url)
print(leaf.wss_url)
print(leaf.discovery_url)
print(leaf.monitor_url)
print(leaf.status().bridge_state)
tls_context = ssl.create_default_context(cafile=str(leaf.ca_cert_file))
# Local TCP listener stays on loopback; WSS and discovery bind to the LAN address by default.
hub = KinopioHub(servers=[leaf.wss_url], tls=tls_context)
start_leaf_node() returns a LeafNodeHandle with status() and stop() plus the stable URLs
client_url, wss_url, discovery_url, and monitor_url. Binary resolution happens in this
order: explicit binary_path, system PATH, an existing user cache entry, then the latest stable
official nats-server release. If you do not provide PEM files, KinopioHub generates a local CA
and a runtime leaf certificate automatically, but it does not modify the system trust store for
you. When backbone_servers is configured but unavailable, the local node still starts and reports
bridge_state == "connecting" while remaining usable for local traffic.
The default cache root is ~/Library/Caches/kinopio-hub on macOS,
%LOCALAPPDATA%\kinopio-hub on Windows, and $XDG_CACHE_HOME/kinopio-hub or
~/.cache/kinopio-hub on Linux. Set KINOPIO_HUB_CACHE_DIR to override it. The default install
already includes the Python dependencies used by the leaf runtime, so there is no separate leaf
extra to install. The local leaf runtime is intended for environments that can spawn a local
nats-server subprocess and open local TCP or UDP sockets; restricted browser and serverless
runtimes are outside this support boundary.
The manual leaf discovery manifest also includes JS-compatible lease metadata such as
expiresAt, leaseExpiresAt, leaderEpoch, isLeader, and candidateRole, so browser-side
consumers can normalize it without requiring the auto-election runtime.
Auto Leaf
enable_auto_leaf() builds on the manual runtime. It uses UDP multicast heartbeats as the base
coordination layer, optionally advertises the leader via mDNS when zeroconf is available, keeps a
stable node_id under the KinopioHub user cache, and exposes a high-level status API for leader
discovery and failover.
import time
from kinopio_hub.leaf import AutoLeafOptions, enable_auto_leaf
handle = enable_auto_leaf(
AutoLeafOptions(
discovery_namespace="studio-a",
backbone_servers=("nats://127.0.0.1:7422",),
leader_missing_grace_ms=10_000,
)
)
try:
while True:
status = handle.status()
print(status.state, status.role, status.current_leader)
time.sleep(1)
finally:
handle.stop()
The public state machine is:
discoveringfollowing-leaderleader-missing-graceelectingstarting-leafleaderstopped
AutoLeafHandle also exposes state(), role(), current_leader(), status(), and stop().
Leader discovery manifests include the JS-facing fields version, expiresAt,
leaderEpoch, advertisedHostname, wssUrl, fallbackServers, bridgeState,
backboneServers, backboneRttMs, discoveryUrl, leaseExpiresAt, nodeId,
discoveryNamespace, isLeader, and candidateRole. A healthier leader will not be preempted immediately; another node must stay at
least 50ms ahead in measured backbone RTT for a sustained window before it can take over.
KinopioHub itself does not implement the browser-side local probe used by the JS package, but it
does emit a compatible discovery manifest that browser clients can consume.
CLI
The kinopio-hub console script exposes the two public leaf workflows:
kinopio-hub leaf start \
--backbone-server nats://127.0.0.1:7422 \
--lan-bind-address 192.168.1.20 \
--json
kinopio-hub leaf auto \
--discovery-namespace studio-a \
--backbone-server nats://127.0.0.1:7422
Both commands support --backbone-server (repeatable), --binary-path, --lan-bind-address,
--client-port, --websocket-port, --discovery-port, --monitor-port, and --json.
kinopio-hub leaf auto also adds --discovery-namespace and --leader-missing-grace-ms.
--lan-bind-address maps to the advertised WSS and discovery address; the local TCP client and
monitor listeners still stay on loopback unless you use the Python API to override them explicitly.
--json emits machine-readable startup or status snapshots for automation and shell tooling.
Serialization
The default serialization rules are:
Nonebecomes an empty payloadbytes,bytearray, andmemoryvieware passed through as bytesstris encoded as UTF-8- everything else is serialized as JSON
If a custom codec is supplied, its encode() and decode() methods are used first.
Service Errors
If a service handler raises an exception, KinopioHub sends a response payload shaped like this:
{"error": True, "message": "details"}
This keeps request/reply flows explicit and observable on the wire.
Examples
See the examples directory for complete runnable programs:
Development
See How_To_Dev.md for development setup and verification steps.
License
GPL-3.0-or-later
Project details
Release history Release notifications | RSS feed
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 kinopio_hub-0.1.1.tar.gz.
File metadata
- Download URL: kinopio_hub-0.1.1.tar.gz
- Upload date:
- Size: 51.6 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.11
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
cdddee5b26a2c72488b162fdbf79cc132f00b4626d8a14aa16f1bf9fd73b23d3
|
|
| MD5 |
da189cc370633016c94663975abd2491
|
|
| BLAKE2b-256 |
efaa9a1ab75e99ac228f2ac7588ee4461f92deb5f10f275a04bb07b0784ae396
|
File details
Details for the file kinopio_hub-0.1.1-py3-none-any.whl.
File metadata
- Download URL: kinopio_hub-0.1.1-py3-none-any.whl
- Upload date:
- Size: 47.3 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.11
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
685fc38638a2008172de8d24c789d816a6049f14e49f8d30e2fcd71c8cedb015
|
|
| MD5 |
a3dc88971cf3b9ae92af47ebfe73e0a0
|
|
| BLAKE2b-256 |
037c24af1c2a10d30d8a2422027e8982f64afec6fc1f41c4c897c66f8afa70ec
|