Developer Guide — Client Library¶
Twenty lines of Python and you're holding a conversation with your hive. That's the
promise of hivemind-bus-client (the repo is named hivemind-websocket-client): you
create a client, call connect(), and send a message — and everything gnarly underneath
happens without you lifting a finger. The handshake, the key derivation, the
authenticated encryption, the wire serialization: all handled. You write "what time is
it?"; the library does the cryptography.
In a nutshell
- Talk to hivemind-core from Python with
HiveMessageBusClient— send and receive OVOS messages while the library handles the handshake, encryption, and serialization for you. - Three flavours share one message API: sync (thread-based), async (
asyncio), and HTTP (polling). Only the sync client reconnects on its own. - Higher-level helpers cover the things you'd otherwise hand-roll: blocking request/response, binary audio, topology discovery, CASCADE aggregation, end-to-end INTERCOM, and trusted-key identity.
Building your own client in Python?
This page is for writing a program that talks to a hivemind-core server. To simply use a satellite, see Choosing a Satellite instead.
hivemind-bus-client is available on PyPI and GitHub.
Install¶
Libraries¶
| Library | Language | Notes |
|---|---|---|
| hivemind-bus-client | Python | Primary client; WebSocket + HTTP + async + MQTT |
| HiveMind-js | JavaScript | Browser and Node.js |
| ovos-solver-hivemind-plugin | Python | OVOS solver plugin — chat with a hivemind-core server |
Basic connection¶
Everything starts with two lines: make a client, connect it. If you've already run
hivemind-client set-identity, the client picks up host, port, key, and password from the
identity file on its own — you don't pass a thing:
from hivemind_bus_client.client import HiveMessageBusClient
# Reads host, port, key, password from ~/.config/hivemind/_identity.json
client = HiveMessageBusClient()
client.connect()
Or with explicit parameters:
client = HiveMessageBusClient(
host="192.168.1.10",
port=5678,
key="my-access-key",
password="my-password"
)
client.connect()
Sending utterances¶
Connected. Now say something. Sending an utterance is just wrapping an OVOS Message in a
BUS envelope and emitting it — the same BUS message from the protocol page, built by
hand:
from ovos_bus_client.message import Message
from hivemind_bus_client.message import HiveMessage, HiveMessageType
utt = "What time is it?"
msg = HiveMessage(
HiveMessageType.BUS,
Message("recognizer_loop:utterance", {"utterances": [utt]})
)
client.emit(msg)
Handling responses¶
The reply comes back asynchronously, so you don't wait on the emit — you register a
callback and let it fire when the answer arrives. on_mycroft listens for a specific
inner OVOS message type, and speak is the one that carries spoken text:
def handle_speak(message):
print(f"AI says: {message.data['utterance']}")
client.on_mycroft("speak", handle_speak)
Send, listen, react — that's the entire loop. Everything below is refinement on top of these three moves: blocking instead of callbacks, binary instead of text, many nodes instead of one.
Conversational loop¶
from hivemind_bus_client.client import HiveMessageBusClient
from ovos_bus_client.message import Message
from hivemind_bus_client.message import HiveMessage, HiveMessageType
client = HiveMessageBusClient()
client.connect()
def on_speak(msg):
print(f"Response: {msg.data['utterance']}")
client.on_mycroft("speak", on_speak)
while True:
try:
utterance = input("You: ").strip()
if utterance:
client.emit(HiveMessage(
HiveMessageType.BUS,
Message("recognizer_loop:utterance", {"utterances": [utterance]})
))
except KeyboardInterrupt:
break
The synchronous HiveMessageBusClient runs its WebSocket on a background thread; there is no close() method — interrupting the loop is enough to stop sending. (An explicit close() exists only on the async client in hivemind_bus_client.async_client.)
Sending binary data (audio)¶
HTTP client¶
For environments where a persistent WebSocket is not feasible:
from hivemind_bus_client.http_client import HiveMindHTTPClient
client = HiveMindHTTPClient(host="http://192.168.1.10", port=5679)
client.connect()
Identity management¶
# Write the identity file
hivemind-client set-identity \
--host 192.168.1.10 \
--port 5678 \
--key "my-key" \
--password "my-password" \
--siteid living-room
# Test the stored identity
hivemind-client test-identity
The identity file is at ~/.config/hivemind/_identity.json.
Receiving binary data¶
Inbound binary payloads (TTS audio, files) are dispatched through a
BinaryDataCallbacks instance. Subclass it, override the handlers you care
about, and pass it as bin_callbacks= to the constructor:
from hivemind_bus_client.client import HiveMessageBusClient, BinaryDataCallbacks
class MyBinaryHandler(BinaryDataCallbacks):
def handle_receive_tts(self, bin_data: bytes,
utterance: str, lang: str, file_name: str):
# synthesized TTS audio to play back
with open(file_name, "wb") as f:
f.write(bin_data)
def handle_receive_file(self, bin_data: bytes, file_name: str):
# arbitrary file pushed from hivemind-core
with open(file_name, "wb") as f:
f.write(bin_data)
client = HiveMessageBusClient(bin_callbacks=MyBinaryHandler())
client.connect()
Serialization and encryption¶
You may have noticed you never once touched a cipher. That's deliberate — client.emit()
and client.on() do the whole cryptographic dance for you, in this order, on every
message:
- The
HiveMessageis serialized (JSON or binary framing for protocol v1) - The payload is compressed with zlib if enabled
- The result is encrypted with AES-256-GCM using the session key
- The encrypted payload is encoded (Z85 + Base91) for text transport
Decorator helpers¶
hivemind_bus_client.decorators provides decorators that register a function as
a handler for a given message type. Each takes a bus= argument (the connected
client). The OVOS-bus variant, on_mycroft_message, filters on the inner OVOS
message type:
from hivemind_bus_client.client import HiveMessageBusClient
from hivemind_bus_client.decorators import on_mycroft_message
client = HiveMessageBusClient()
client.connect()
@on_mycroft_message(payload_type="speak", bus=client)
def on_speak(message):
print(f"AI says: {message.data['utterance']}")
Other decorators target specific HiveMessageType envelopes:
on_hive_message, on_payload, on_ping, on_broadcast, on_propagate,
on_escalate, on_handshake, on_hello, on_query, on_cascade,
on_rendezvous, on_third_party, and on_shared_bus.
Topology mapping (PING flood)¶
HiveMind discovers its topology with a PING flood: each node answers a PING by
re-emitting its own PING carrying the same flood_id, propagated across the
hive. There is no separate PONG message type — replies are just more PINGs.
import uuid, time
from hivemind_bus_client.message import HiveMessage, HiveMessageType
flood_id = str(uuid.uuid4())
ping_inner = HiveMessage(
HiveMessageType.PING,
payload={
"flood_id": flood_id,
"timestamp": time.time(),
"peer": f"{client.site_id}::{client.session_id}",
"site_id": client.site_id,
}
)
# PING must always be wrapped in PROPAGATE so it floods the hive
ping_outer = HiveMessage(HiveMessageType.PROPAGATE, payload=ping_inner)
client.emit(ping_outer)
def on_ping_reply(message):
payload = message.payload
route = message.route # list of {source, targets} hop records
print(f"PING from {payload['peer']} via {len(route)} hops")
# Responses arrive as PING messages (re-emitted with the same flood_id)
client.on(HiveMessageType.PING, on_ping_reply)
For automated topology collection use HiveMapper from
hivemind_bus_client.hive_map. Call mapper.start_ping(flood_id), feed each
received inner PING to mapper.on_ping(ping_msg), and read back the discovered
mapper.nodes / mapper.edges.
Request / response¶
Callbacks are great for a UI, awkward for a script that just wants an answer now. For
those, the client ships blocking helpers: emit a message, then wait on the background
thread until the matching reply lands (or a timeout gives up). All take a timeout (seconds, default
3.0) and return the received message or None on timeout.
def wait_for_message(self, message_type, timeout=3.0)
def wait_for_payload(self, payload_type,
message_type=HiveMessageType.BUS, timeout=3.0)
def wait_for_mycroft(self, mycroft_msg_type, timeout=3.0)
def wait_for_response(self, message, reply_type=None, timeout=3.0)
def wait_for_payload_response(self, message, payload_type,
reply_type=None, timeout=3.0)
wait_for_message— wait for the nextHiveMessageof a givenHiveMessageType.wait_for_payload— wait for aHiveMessageofmessage_typewhose inner payload type matchespayload_type(defaults the envelope toBUS).wait_for_mycroft— convenience wrapper:wait_for_payload(mycroft_msg_type, message_type=HiveMessageType.BUS). Use it to wait for an inner OVOS message type (e.g."speak").wait_for_response— emitmessage, then wait for a reply.reply_typedefaults tomessage.msg_type. A MycroftMessagewaits on the inner payload type; aHiveMessagewaits on the envelope type.wait_for_payload_response— emitmessage, then wait for a reply envelope ofreply_typewhose inner payload matchespayload_type.
The request/response idiom (ask once, block for the answer):
from ovos_bus_client.message import Message
from hivemind_bus_client.message import HiveMessageType
reply = client.wait_for_payload_response(
message=Message("recognizer_loop:utterance",
{"utterances": ["what time is it?"]}),
payload_type="speak",
reply_type=HiveMessageType.BUS,
timeout=5,
)
if reply is not None:
print(reply.payload.data["utterance"])
The same five helpers exist on the async client as coroutines — await
client.wait_for_response(...).
INTERCOM — end-to-end satellite-to-satellite¶
Everything so far has been you talking to hivemind-core. emit_intercom is different: it
lets you whisper to one specific other peer, sealed so tightly that the servers relaying
the message can't read a word of it. Reach for it when two devices need a private side
channel — one satellite nudging another — and hivemind-core should only ferry the envelope,
never open it. The one thing you must have in hand is the recipient's RSA public key:
Notice you never manage a session key here — the library builds a fresh one per message.
On the WebSocket/async client the payload is hybrid-encrypted via
hybrid_encrypt (encryption.py): a random AES-256 key encrypts the serialized
message with AES-GCM, that AES key is RSA-OAEP-wrapped with pubkey, and the
ciphertext is signed with this node's private key. The envelope is a dict with
base64 fields encrypted_key, ciphertext, tag, nonce, and signature,
wrapped in a HiveMessageType.INTERCOM message.
recipient_pubkey = client.identity.trusted_keys["living-room-satellite"]
client.emit_intercom(
HiveMessage(HiveMessageType.BUS,
Message("speak", {"utterance": "psst"})),
pubkey=recipient_pubkey,
)
For the inbound side to accept and inject an untargeted INTERCOM message, the
sender's key must be in the receiver's trusted_keys (see Identity & trusted
keys below); a message whose target_public_key matches the receiver is always
accepted.
HTTP wire format differs.
HiveMindHTTPClient.emit_intercomdoes not use the hybrid envelope. It RSA-encrypts the whole serialized message withencrypt_RSA(pubkey, ...), signs it, and sends a payload of just{"ciphertext": ..., "signature": ...}(both base64). Don't mix the two formats across transports.
Async client¶
AsyncHiveMessageBusClient (hivemind_bus_client.async_client) mirrors the sync
client on asyncio + the websockets library. Install with the extra:
The connect lifecycle is explicit and all I/O is awaited:
import asyncio
from hivemind_bus_client.async_client import AsyncHiveMessageBusClient
from hivemind_bus_client.message import HiveMessage, HiveMessageType
from ovos_bus_client.message import Message
async def main():
bus = AsyncHiveMessageBusClient(key="my-key", password="my-password",
host="192.168.1.10", port=5678)
await bus.connect() # opens WS, binds protocol, awaits handshake
try:
await bus.emit(HiveMessage(
HiveMessageType.BUS,
Message("recognizer_loop:utterance", {"utterances": ["hi"]})))
reply = await bus.wait_for_mycroft("speak", timeout=5)
if reply is not None:
print(reply.payload.data["utterance"])
finally:
await bus.close() # closes WS, cancels the receive task
asyncio.run(main())
Key coroutines: connect(bus=None, protocol=None, site_id=None), close(),
emit(...), emit_mycroft(...), emit_intercom(...), and the five wait_for_*
waiters. Handler registration is synchronous (so existing protocol handlers
work unchanged): on(event, func), once(event, func) (fire-once), and
remove(event, func).
Use the async client when your app already runs an event loop (aiohttp/FastAPI,
discord.py, etc.). Use the sync HiveMessageBusClient for scripts and
thread-based apps — it runs its WebSocket on a background thread and needs no loop.
HTTP client specifics¶
HiveMindHTTPClient (hivemind_bus_client.http_client) is a
threading.Thread. Its __init__ calls self.start(), so the polling thread
begins as soon as you construct it — you still must call connect() before
emitting (it raises ConnectionAbortedError otherwise).
from hivemind_bus_client.http_client import HiveMindHTTPClient
client = HiveMindHTTPClient(host="http://192.168.1.10", port=5679)
client.connect() # POST /connect, then awaits handshake
resp = client.emit(some_hive_message) # returns a requests.Response
emit(message, binary_type=...)returns therequests.ResponsefromPOST {base_url}/send_message(form fieldmessage), notNone.- The background
run()loop pollsGET {base_url}/get_messagesandGET {base_url}/get_binary_messagesonce per second and feeds each result toon_message. There is no server push; latency is bounded by that poll interval. - Register handlers with
on(hive_type, func)for envelopes andon_mycroft(ovos_type, func)for inner OVOS messages;remove/remove_mycroftunregister them. disconnect()(POST/disconnect) andshutdown()stop the loop.
CASCADE aggregation¶
A CASCADE query is answered by many nodes across the hive. CascadeAggregator
(hivemind_bus_client.protocol) collects those responses over a window and picks
one winner.
class CascadeAggregator:
def __init__(self, timeout, select_callback, emit_callback,
expected_responses=None): ...
def add_response(self, message): ...
def cancel(self): ...
When the first response arrives a timeout-second timer starts; subsequent
responses are buffered. It resolves when the timer fires or once
expected_responses responses have arrived (whichever is first), at which point
select_callback(List[HiveMessage]) -> Optional[HiveMessage] chooses the winner
and emit_callback(HiveMessage) delivers it. cancel() drops the timer and
buffered responses without emitting.
Inside HiveMindSlaveProtocol.handle_cascade, a node builds one per query:
select_callback defaults to random.choice (or the protocol's
cascade_select_callback), expected_responses defaults to the number of nodes
known to the HiveMapper, and emit_callback injects the winner's inner BUS
payload onto the internal bus.
Identity & trusted keys¶
NodeIdentity (hivemind_bus_client.identity) wraps the identity file
(~/.config/hivemind/_identity.json via JsonConfigXDG, or pass
identity_file=). Beyond the connection fields (name, access_key,
password, default_master, default_port, site_id, public_key,
private_key):
identity = NodeIdentity()
identity.create_keys() # generate + store an RSA keypair
identity.add_trusted_key("living-room-satellite", # alias -> public key
peer_public_key) # True if added, False if alias exists
identity.save() # persist to the identity file
identity.reload() # re-read from disk
print(identity.trusted_keys) # {alias: pubkey, ...}
trusted_keys is the alias→public-key mapping used to verify peers in PROPAGATE,
CASCADE, and INTERCOM handling — only messages from a trusted key (or explicitly
targeted at this node) are injected onto the bus. Related helpers:
remove_trusted_key(alias), is_trusted_key(pubkey),
get_trusted_alias(pubkey).
Reconnection¶
Networks drop. Whether your client copes on its own depends entirely on which one you picked — and this catches people out, so it's worth stating plainly before you ship.
The sync HiveMessageBusClient reconnects automatically; the async and HTTP
clients do not. On a dropped websocket the sync client's on_error clears its
crypto/handshake state and then falls through to the inherited OVOS bus-client
reconnect loop (exponential backoff capped at 60 s, reset to 5 s on each
successful connect); it re-runs the handshake after reconnecting. You do not need
your own supervisor for the sync client.
The async (AsyncHiveMessageBusClient) and HTTP (HiveMindHTTPClient)
clients have no built-in retry/backoff loop. If your deployment uses either
and needs resilience against dropped connections, build your own supervisor:
detect the drop (e.g. emit failures, the async receive loop emitting "close",
or your own heartbeat) and re-run connect() on a fresh client with backoff.
Encodings & ciphers¶
The JSON transport encoding and the symmetric cipher are configurable. The enums
live in hivemind_bus_client.encryption:
SupportedEncodings (string values): JSON_B91 ("JSON-B91"), JSON_Z85B
("JSON-Z85B"), JSON_Z85P ("JSON-Z85P"), JSON_B64 ("JSON-B64"),
JSON_URLSAFE_B64 ("JSON-URLSAFE-B64"), JSON_B32 ("JSON-B32"), JSON_HEX
("JSON-HEX").
SupportedCiphers: AES_GCM ("AES-GCM"), CHACHA20_POLY1305
("CHACHA20-POLY1305").
Clients default to JSON_HEX + AES_GCM; the actual values in use are negotiated
during the handshake. To force them, set the attributes after construction:
from hivemind_bus_client.encryption import SupportedEncodings, SupportedCiphers
client = HiveMessageBusClient()
client.json_encoding = SupportedEncodings.JSON_B64
client.cipher = SupportedCiphers.CHACHA20_POLY1305
client.connect()
See also¶
- Protocol Specification — wire format, message types, routing
- Testing Guide — writing tests with the in-process harness
- CLI Reference —
hivemind-clientcommands
Source¶
Validated against the HiveMind source:
hivemind_bus_client/client.py—HiveMessageBusClient,BinaryDataCallbacks, thewait_for_*helpers,emit_intercom, reconnection, and the encoding/cipher attributeshivemind_bus_client/encryption.py—hybrid_encrypt,SupportedEncodings,SupportedCiphershivemind_bus_client/identity.py—NodeIdentity,add_trusted_key,trusted_keys