Skip to main content

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 methodParametersReturnsBehavior and errors
StreamWriter.write(data, topic=None)Any serializer-supported value and an optional topicTransport-specific result or NoneQueues the message when queue_size > 0; a full queue drops its oldest item. A closed writer drops new writes.
StreamWriter.close()—NoneDrains 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()—NoneStops background reading and closes transport resources. Safe to call repeatedly.
RpcRequester.call(request, timeout=None)Raw request value and timeout in secondsDecoded responseRaises 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 deadlineUnwrapped method resultRaises JsonRpcError for schema errors and RuntimeError after close. Dynamic calls such as client.move(x=1) use the same path.
RpcRequester.close()—NoneReleases requester resources.
RpcResponder.respond(handler=None, timeout=None)handler(request) -> response; handler is optional when a schema is attachedTrue after one request, False on timeoutReceives one request, dispatches it, and replies. Transport or handler errors propagate.
RpcResponder.close()—NoneReleases 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​

ClassConstructor and important options
ZmqStreamWriterZmqStreamWriter(endpoint, serializer=None, queue_size=10, bind=True, delivery="reliable")
ZmqStreamReaderZmqStreamReader(endpoint, topic="", serializer=None, queue_size=10, bind=False, delivery="reliable"); topic may be a string or list
ZMQRpcRequesterZMQRpcRequester(endpoint, serializer=None, name=None, identity=None, ack_timeout=2.0, schema=None)
ZMQRpcResponderZMQRpcResponder(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​

ClassConstructor and public methods
MqttConnectionMqttConnection(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.

APIPurpose
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_idsReturns this participant's ID or the currently connected remote IDs.
is_connectedReports 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 methodPurpose
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).

FamilyClasses and helpers
PrimitiveBoolFrame, IntFrame, FloatFrame, StringFrame, BytesFrame, ListFrame, DictFrame
ImageImageFrameRaw; ImageFrameCV.from_cv_image() / to_cv_image(); ImageFrameJpeg.from_np_image() / to_np_image()
AudioAudioFrameRaw; 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​

APIRole
BaseNodeThreaded lifecycle with setup(), process(), cleanup(), pause(), resume(), terminate(), paused(), and terminating().
SourceNodeRuns a producer around a StreamWriter.
SinkNodeRuns a consumer around a StreamReader.
ProcessNodeCombines stream input and output.
ServerNodeRuns an RpcResponder; accepts a handler or a responder-attached schema.
ZconfDiscoveryadvertise_node(), resolve_node(), list_nodes(), pick_best_ip(), stop_advertising(), and close() for Zeroconf discovery.
McastDiscoveryadvertise(), scan(), and stop_advertising() for multicast discovery.

Exceptions and lifecycle​

  • Catch AckTimeoutError separately from ReplyTimeoutError when the distinction affects retry policy.
  • Catch JsonRpcError for 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.