Python API
Package: luxai-magpie · Repository: luxai-qtrobot/magpie · Python 3.9+
Language references: Python · C++ · TypeScript / JavaScript
Core communication API
These four classes are the stable application boundary. Concrete ZeroMQ, MQTT, and WebRTC classes preserve the same method semantics.
| Class and method | Parameters | Returns | Behavior and errors |
|---|---|---|---|
StreamWriter.write(data, topic=None) | Any serializer-supported value and an optional topic | Transport-specific result or None | Queues the message when queue_size > 0; a full queue drops its oldest item. A closed writer drops new writes. |
StreamWriter.close() | — | None | Drains pending writes, closes transport resources, and is safe to call repeatedly. |
StreamReader.read(timeout=None) | Timeout in seconds; None waits indefinitely | (data, topic) | Raises TimeoutError when no item arrives before the deadline. Returns None if closed while waiting. |
StreamReader.close() | — | None | Stops background reading and closes transport resources. Safe to call repeatedly. |
RpcRequester.call(request, timeout=None) | Raw request value and timeout in seconds | Decoded response | Raises AckTimeoutError when no responder acknowledges, or ReplyTimeoutError when an acknowledged request does not finish in time. |
RpcRequester.call(method, **params) | Method name and named parameters when a schema is attached; use _timeout= for the deadline | Unwrapped method result | Raises JsonRpcError for schema errors and RuntimeError after close. Dynamic calls such as client.move(x=1) use the same path. |
RpcRequester.close() | — | None | Releases requester resources. |
RpcResponder.respond(handler=None, timeout=None) | handler(request) -> response; handler is optional when a schema is attached | True after one request, False on timeout | Receives one request, dispatches it, and replies. Transport or handler errors propagate. |
RpcResponder.close() | — | None | Releases responder resources. |
handle_once() remains as a backward-compatible alias for respond(). New code should use respond().
Streaming example
from luxai.magpie.transport import ZmqStreamReader, ZmqStreamWriter
writer = ZmqStreamWriter("tcp://*:5555", queue_size=10)
reader = ZmqStreamReader("tcp://127.0.0.1:5555", topic="robot/state")
try:
writer.write({"battery": 0.82}, topic="robot/state")
value, topic = reader.read(timeout=5.0)
print(topic, value)
finally:
reader.close()
writer.close()
RPC example
from luxai.magpie.transport import ZMQRpcResponder
server = ZMQRpcResponder("tcp://*:5556")
def status(request):
return {"ready": True, "echo": request}
try:
while True:
server.respond(status, timeout=1.0)
finally:
server.close()
From a separate client process:
from luxai.magpie.transport import ZMQRpcRequester
client = ZMQRpcRequester("tcp://127.0.0.1:5556")
try:
reply = client.call({"action": "status"}, timeout=5.0)
print(reply)
finally:
client.close()
Transport classes
ZeroMQ
| Class | Constructor and important options |
|---|---|
ZmqStreamWriter | ZmqStreamWriter(endpoint, serializer=None, queue_size=10, bind=True, delivery="reliable") |
ZmqStreamReader | ZmqStreamReader(endpoint, topic="", serializer=None, queue_size=10, bind=False, delivery="reliable"); topic may be a string or list |
ZMQRpcRequester | ZMQRpcRequester(endpoint, serializer=None, name=None, identity=None, ack_timeout=2.0, schema=None) |
ZMQRpcResponder | ZMQRpcResponder(endpoint, serializer=None, name=None, bind=True, schema=None) |
bind chooses whether the class owns the listening endpoint. delivery="latest" favors current stream state; "reliable" uses normal queued delivery. ZmqStreamWriter.wait_connect(timeout) can be used where startup ordering matters.
MQTT
| Class | Constructor and public methods |
|---|---|
MqttConnection | MqttConnection(uri, client_id=None, protocol_version=5, keepalive=60, options=None) |
MqttConnection.connect(timeout=10.0) | Connects and returns connection success. Supported URI schemes include mqtt://, mqtts://, ws://, and wss://. |
MqttConnection.disconnect() | Stops reconnect handling and closes the broker connection. |
MqttConnection.publish(topic, payload, qos=None, retain=None) | Low-level byte publication; stream writers normally call this for you. |
MqttConnection.add_subscription(...) / remove_subscription(...) | Registers or removes low-level callbacks. |
MqttConnection.is_connected() | Returns current broker connection state. |
MqttStreamWriter | (connection, serializer=None, queue_size=10, qos=None, retain=None) |
MqttStreamReader | (connection, topic, serializer=None, queue_size=10, qos=None); MQTT wildcards are supported. |
MqttRpcRequester | (connection, service_name, serializer=None, name=None, ack_timeout=2.0, qos=None, schema=None) |
MqttRpcResponder | (connection, service_name, serializer=None, name=None, qos=None, schema=None) |
One connection can be shared by multiple readers, writers, requesters, and responders. Close those children before calling disconnect().
from luxai.magpie.transport import MqttConnection, MqttStreamWriter
connection = MqttConnection("mqtts://broker.example.com:8883", options=mqtt_options)
connection.connect(timeout=10)
writer = MqttStreamWriter(connection, qos=1, retain=False)
try:
writer.write({"online": True}, "robot/01/status")
finally:
writer.close()
connection.disconnect()
MqttOptions groups MqttTlsOptions, MqttAuthOptions, MqttSessionOptions, MqttReconnectOptions, MqttWillOptions, and MqttDefaultsOptions so TLS, credentials, session behavior, reconnect policy, last will, and default QoS/retain settings stay explicit.
WebRTC
Create a shared WebRTCConnection with a signaling factory, call connect(), then construct stream or RPC children around it. A connection manages one WebRTC link per remote peer.
| API | Purpose |
|---|---|
WebRTCConnection.with_zmq(endpoint, session_id, bind=..., options=..., role="mesh", multiplex=False) | Rendezvous on a local/LAN ZeroMQ endpoint. Enable multiplex on every participant for multiple peers. |
WebRTCConnection.with_mqtt(broker_uri, session_id, options=..., reconnect=..., role="mesh") | Rendezvous through MQTT, including peers behind NAT. |
WebRTCConnection.with_http(signal_url, session_id, headers=..., options=..., reconnect=..., role="mesh") | Rendezvous through an HTTP signaling service. |
connect(timeout=None) / disconnect() | Starts signaling, waits for the first peer, or closes all peer links. |
peer_id / peer_ids | Returns this participant's ID or the currently connected remote IDs. |
is_connected | Reports whether at least one remote peer is connected. |
WebRtcStreamWriter(connection, queue_size=10) | Sends values and frames. Media frames may use negotiated media tracks. |
WebRtcStreamReader(connection, topic, queue_size=10) | Reads data or media frames for a topic. |
WebRTCRpcRequester(connection, service_name, name=None, ack_timeout=2.0, schema=None) | Calls a peer service. |
WebRTCRpcResponder(connection, service_name, name=None, schema=None) | Responds for a peer service. |
role accepts "mesh" (default, connect every participant), "host" (accept clients), or "client" (connect to hosts, not other clients). Stream writes go to all connected peers; RPC replies are routed back to the requesting peer.
WebRTCOptions configures STUN/TURN, codecs, media-channel use, and connection behavior. WebRTCTurnServer describes a TURN URL and credentials.
Schemas
| Class or method | Purpose |
|---|---|
BaseSchema.dispatch(request) | Converts one protocol request into a response. |
BaseSchema.wrap(method, params=None) | Creates the requester-side envelope. |
BaseSchema.unwrap(response) | Returns a result or raises the protocol error. |
JsonRpcSchema.register(name, func=None, description=None, input_schema=None, output_schema=None) | Registers a JSON-RPC method and optional JSON Schemas. |
JsonRpcSchema.method(name=None) | Decorator that registers a Python function and derives useful metadata. |
JsonRpcSchema.handler(name) | Attaches a handler to a method loaded from a contract. |
from_json(...), from_json_string(...), from_json_file(...) | Loads method definitions from JSON data. |
McpSchema(name="magpie", version="1.0.0") | Adds MCP initialization, tool discovery, calls, and ping behavior. |
McpSchema.register(...) | Registers a method and exposes it as an MCP tool. |
from luxai.magpie.schema import JsonRpcSchema
from luxai.magpie.transport import ZMQRpcResponder
schema = JsonRpcSchema()
@schema.method()
def add(a: float, b: float) -> float:
"""Add two numbers."""
return a + b
server = ZMQRpcResponder("tcp://*:5556", schema=schema)
while True:
server.respond(timeout=1.0)
Frames
Every frame carries gid, id, name, and timestamp. Send frame.to_dict() over a general stream and reconstruct it with Frame.from_dict(data).
| Family | Classes and helpers |
|---|---|
| Primitive | BoolFrame, IntFrame, FloatFrame, StringFrame, BytesFrame, ListFrame, DictFrame |
| Image | ImageFrameRaw; ImageFrameCV.from_cv_image() / to_cv_image(); ImageFrameJpeg.from_np_image() / to_np_image() |
| Audio | AudioFrameRaw; AudioFrameFlac.from_pcm() / to_pcm() |
from luxai.magpie.frames import DictFrame, Frame
outgoing = DictFrame(value={"battery": 0.82})
wire_value = outgoing.to_dict()
incoming = Frame.from_dict(wire_value)
Custom frame subclasses are registered automatically. For cross-language use, give every implementation the same class name and wire fields.
Serialization
MsgpackSerializer is the default, not a requirement. A serializer implements:
class BaseSerializer:
def serialize(self, data: object) -> bytes: ...
def deserialize(self, byte_data: bytes) -> object: ...
Pass a serializer instance to compatible transport constructors. Every sender and receiver on that path must agree on the format. See Extending MAGPIE for a JSON example.
Nodes and discovery
| API | Role |
|---|---|
BaseNode | Threaded lifecycle with setup(), process(), cleanup(), pause(), resume(), terminate(), paused(), and terminating(). |
SourceNode | Runs a producer around a StreamWriter. |
SinkNode | Runs a consumer around a StreamReader. |
ProcessNode | Combines stream input and output. |
ServerNode | Runs an RpcResponder; accepts a handler or a responder-attached schema. |
ZconfDiscovery | advertise_node(), resolve_node(), list_nodes(), pick_best_ip(), stop_advertising(), and close() for Zeroconf discovery. |
McastDiscovery | advertise(), scan(), and stop_advertising() for multicast discovery. |
Exceptions and lifecycle
- Catch
AckTimeoutErrorseparately fromReplyTimeoutErrorwhen the distinction affects retry policy. - Catch
JsonRpcErrorfor method-not-found, invalid-parameter, and server-side schema errors. - A timeout does not prove a mutating RPC was never executed. Use idempotency keys before automatic retries.
- Use
try/finally, close child components first, and disconnect their shared MQTT or WebRTC connection last.
The Python source tree is the canonical signature reference. Complete programs are available in the Python examples.