Extending MAGPIE
MAGPIE is designed for extension at clear boundaries. New transports implement I/O; serializers encode values; frames define portable data; schemas define service methods.
Add a streaming transport
In Python, subclass the abstract writer and reader. Shared queueing, lifecycle, and the public API remain in the base classes.
from luxai.magpie.transport import StreamReader, StreamWriter
class MyStreamWriter(StreamWriter):
def __init__(self, endpoint, serializer, queue_size=10):
self.endpoint = endpoint
self.serializer = serializer
self.client = open_my_transport(endpoint)
super().__init__(name="MyStreamWriter", queue_size=queue_size)
def _transport_write(self, value, topic):
self.client.send(topic, self.serializer.serialize(value))
def _transport_close(self):
self.client.close()
class MyStreamReader(StreamReader):
def __init__(self, endpoint, topic, serializer, queue_size=10):
self.serializer = serializer
self.client = open_my_subscription(endpoint, topic)
super().__init__(name="MyStreamReader", queue_size=queue_size)
def _transport_read_blocking(self, timeout=None):
payload, topic = self.client.receive(timeout=timeout)
return self.serializer.deserialize(payload), topic
def _transport_close(self):
self.client.close()
Preserve the meaning of write, read, timeouts, bounded queues, and idempotent close. Keep transport configuration in the concrete class rather than adding transport branches to application code.
Add an RPC transport
Implement request/reply I/O behind RpcRequester and RpcResponder:
class MyRpcRequester(RpcRequester):
def _transport_call(self, request, timeout=None):
return self.client.call(request, timeout=timeout)
def _transport_close(self):
self.client.close()
class MyRpcResponder(RpcResponder):
def _transport_recv(self, timeout=None):
return self.server.receive(timeout=timeout) # request, client context
def _transport_send(self, response, client_context):
self.server.reply(client_context, response)
def _transport_close(self):
self.server.close()
Document correlation, acknowledgement, timeout, cancellation, and duplicate-delivery behavior. These semantics matter more than the socket mechanics.
Add a serializer
import json
from luxai.magpie.serializer import BaseSerializer
class JsonSerializer(BaseSerializer):
def serialize(self, value):
return json.dumps(value).encode("utf-8")
def deserialize(self, payload):
return json.loads(payload.decode("utf-8"))
Pass the serializer when constructing compatible readers and writers. All endpoints on that path must use the same format.
Add a frame type
from dataclasses import dataclass, field
from luxai.magpie.frames import Frame
@dataclass
class LidarFrame(Frame):
points: list = field(default_factory=list)
coordinate_frame: str = "sensor"
def __post_init__(self):
super().__post_init__()
The Python frame registry records subclasses automatically. For cross-language use, implement the same frame name and wire fields in each language and add round-trip compatibility tests.
Extension checklist
- Keep the existing public abstraction intact.
- Define ownership and shutdown behavior.
- Bound memory, queues, retries, and network waits.
- Separate transport errors from schema/application errors.
- Add unit tests plus at least one end-to-end example.
- Add cross-language fixtures for wire-visible changes.
- Document optional dependencies and security assumptions.