TypeScript API
Package: @luxai-qtrobot/magpie · Repository: luxai-qtrobot/magpie-js · Environments: browser and Node.js
Language references: Python · C++ · TypeScript / JavaScript
Runtime support
| Capability | Node.js | Browser |
|---|---|---|
| MQTT over TCP/TLS | ✓ | — |
| MQTT over WebSocket | ✓ | ✓ |
| WebRTC | — | ✓ |
| ESM and CommonJS | ✓ | ✓ with a bundler |
| UMD/CDN bundle | — | ✓ |
| TypeScript declarations | ✓ | ✓ |
WebRTC uses native browser APIs. Browser MQTT connections must use ws:// or wss://; prefer wss:// outside local development.
Core communication API
Network operations are asynchronous. Await connection setup, reads, writes, calls, and disconnection.
| Class and method | Parameters and return | Behavior |
|---|---|---|
StreamWriter.write(data, topic) | unknown, topic string → Promise<void> | Serializes and publishes one value. |
StreamWriter.close() | void | Releases writer resources; it does not disconnect a shared connection. |
StreamReader.read(timeout?) | Seconds → Promise<[data, topic]> | Resolves with one item or rejects with TimeoutError. |
StreamReader.close() | void | Unsubscribes, rejects pending readers, and releases local resources. |
RpcRequester.call(request, timeout?) | Raw request → Promise<unknown> | Rejects with AckTimeoutError or ReplyTimeoutError at the corresponding RPC stage. |
RpcRequester.call(method, params?, timeout?) | Method name and parameter object when a schema is configured | Wraps and unwraps JSON-RPC automatically. A proxy call such as client.add({a: 1, b: 2}) uses the same path. |
RpcRequester.close() | void | Releases requester state and pending calls. |
RpcResponder.respond(handler) | Sync or async (request) => response; returns void | Registers the handler and returns immediately. Incoming requests are dispatched by the event loop. |
RpcResponder.close() | void | Stops handling the service and releases subscriptions/callbacks. |
The former onRequest() name is retained only for backward compatibility. New code should use respond().
Streaming example
import {
MqttConnection,
MqttStreamReader,
MqttStreamWriter,
TimeoutError,
} from '@luxai-qtrobot/magpie'
const connection = new MqttConnection('wss://broker.example.com/mqtt')
await connection.connect()
const writer = new MqttStreamWriter(connection)
const reader = new MqttStreamReader(connection, {topic: 'robot/state'})
try {
await writer.write({battery: 0.82}, 'robot/state')
const [value, topic] = await reader.read(5)
console.log(topic, value)
} catch (error) {
if (!(error instanceof TimeoutError)) throw error
} finally {
reader.close()
writer.close()
await connection.disconnect()
}
RPC example
import {
AckTimeoutError,
MqttConnection,
MqttRpcRequester,
MqttRpcResponder,
ReplyTimeoutError,
} from '@luxai-qtrobot/magpie'
const connection = new MqttConnection('mqtt://broker.example.com:1883')
await connection.connect()
const server = new MqttRpcResponder(connection, 'robot/actions')
server.respond((request) => ({status: 'ok', echo: request}))
const client = new MqttRpcRequester(connection, 'robot/actions')
try {
const response = await client.call({action: 'status'}, 5)
console.log(response)
} catch (error) {
if (error instanceof AckTimeoutError) console.error('No responder acknowledged the request')
else if (error instanceof ReplyTimeoutError) console.error('The responder did not reply in time')
else throw error
}
MQTT API
MqttConnection
new MqttConnection(uri, options?: MqttOptions & {clientId?: string})
| Method | Purpose |
|---|---|
connect(timeout=10_000) | Connects to the broker; timeout is in milliseconds. |
disconnect() | Gracefully disconnects and returns Promise<void>. |
publish(topic, payload, qos?, retain?) | Publishes raw Uint8Array data. Stream writers normally call this for you. |
addSubscription(topic, callback, qos?) | Registers a callback. MQTT + and # wildcards are supported. |
removeSubscription(topic, callback) | Removes a callback and unsubscribes when no listeners remain. |
MqttOptions groups authentication, last-will, session, reconnect, and default QoS/retain settings. Use TLS options through the selected secure URI and runtime configuration.
MQTT streams and RPC
| Class | Constructor |
|---|---|
MqttStreamWriter | (connection, {serializer?, qos?, retain?}?) |
MqttStreamReader | (connection, {topic, serializer?, qos?, queueSize?}) |
MqttRpcRequester | (connection, serviceName, {schema?, serializer?, ackTimeout?, qos?, name?}?) |
MqttRpcResponder | (connection, serviceName, {serializer?, qos?, schema?}?) |
One connection can be shared by many components. Close the children before awaiting connection.disconnect().
WebRTC API
WebRTC is browser-only and uses a signaling channel to establish a separate WebRTC link to each remote peer. Application data then travels through those links.
| Factory or method | Purpose |
|---|---|
WebRtcConnection.withMqtt(brokerUrl, sessionId, options?) | Creates a connection using MQTT-over-WebSocket signaling. |
WebRtcConnection.withHttp(signalUrl, sessionId, options?) | Creates a connection using an HTTP signaling relay. Supports static or refreshed request headers. |
connect(timeout?) | Establishes the peer connection; timeout is in seconds. |
disconnect() | Closes peer and signaling resources. |
peerIds | Lists the currently connected remote peer IDs. |
isVideoNegotiated(topic) / isAudioNegotiated(topic) | Reports whether a media topic has a negotiated native track. |
sendVideoTrack(track, topic) / sendAudioTrack(track, topic) | Publishes a browser MediaStreamTrack. |
receiveVideoTrack(topic, peerId?) / receiveAudioTrack(topic, peerId?) | Resolves with a remote media track, optionally from a specific peer. |
The connection option 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 return to the requester.
WebRtcOptions configures STUN servers, TURN servers, ICE policy, data-channel ordering and retransmits, media-channel use, JPEG quality, and negotiated media topics.
| Class | Constructor |
|---|---|
WebRtcStreamWriter | (connection) |
WebRtcStreamReader | (connection, topic, queueSize=10) |
WebRtcRpcRequester | (connection, serviceName, {schema?, ackTimeout?}?) |
WebRtcRpcResponder | (connection, serviceName, {schema?}?) |
const connection = await WebRtcConnection.withHttp(
'https://signal.example.com/webrtc',
'robot-01',
{
headersProvider: async () => ({Authorization: `Bearer ${await getToken()}`}),
role: 'client',
reconnect: true,
webrtcOptions: {videoTopics: ['/camera/color/image']},
},
)
await connection.connect(60)
Schemas
| Class or method | Purpose |
|---|---|
BaseSchema.dispatch(request) | Asynchronously dispatches a decoded request. |
JsonRpcSchema.register(name, handler?, options?) | Registers a method with description and optional input/output JSON Schemas. |
JsonRpcSchema.handler(name, handler) | Attaches an implementation to a method loaded from a contract. |
JsonRpcSchema.fromJSON(definitions) | Builds a schema from shared method definitions. |
wrap(method, params?) | Creates a JSON-RPC request envelope. |
unwrap(response) | Returns a result or throws JsonRpcError. |
McpSchema({name?, version?}) | Adds MCP initialization, tool listing, calls, and ping behavior. |
McpSchema.register(...) | Registers a method and exposes it as an MCP tool. |
import {JsonRpcSchema, MqttRpcResponder} from '@luxai-qtrobot/magpie'
const schema = new JsonRpcSchema()
schema.register(
'add',
(params: unknown) => {
const {a, b} = params as {a: number; b: number}
return a + b
},
{
description: 'Add two numbers',
inputSchema: {
type: 'object',
properties: {a: {type: 'number'}, b: {type: 'number'}},
required: ['a', 'b'],
},
},
)
const server = new MqttRpcResponder(connection, 'math', {schema})
When a responder has a schema, it dispatches incoming requests automatically; no separate handler registration is needed.
Frames
Every frame carries gid, id, name, and timestamp metadata.
| Family | Classes |
|---|---|
| Primitive | BoolFrame, IntFrame, FloatFrame, StringFrame, BytesFrame, ListFrame, DictFrame |
| Image | ImageFrameRaw, ImageFrameJpeg |
| Audio | AudioFrameRaw, AudioFrameFlac |
| Method | Purpose |
|---|---|
frame.toDict() | Produces a plain wire-format object. |
Frame.fromDict(data) | Reconstructs a registered subclass from its name; unknown names fall back to Frame. |
Frame.register(name, factory) | Registers a custom frame factory. |
import {DictFrame, Frame} from '@luxai-qtrobot/magpie'
const outgoing = new DictFrame({value: {battery: 0.82}})
await writer.write(outgoing.toDict(), 'robot/state')
const [raw] = await reader.read(5)
const incoming = Frame.fromDict(raw as Record<string, unknown>)
Serialization
MsgpackSerializer is the default, not a requirement. A custom serializer extends BaseSerializer and implements:
serialize(data: unknown): Uint8Array
deserialize(data: Uint8Array | ArrayBuffer): unknown
Pass it through the serializer option of matching MQTT stream or RPC components. Every endpoint on that path must agree on the format.
MCP adapter
McpTransport connects the official TypeScript MCP SDK to any MAGPIE RpcRequester.
| Method | Purpose |
|---|---|
new McpTransport(requester, timeout=30) | Borrows a requester and sets the call timeout. |
start() | Starts the MCP transport and invokes its startup callback. |
send(message) | Sends an MCP request through MAGPIE RPC. |
close() | Closes adapter state; close the borrowed requester separately. |
Close resources in this order: MCP client, McpTransport, requester, then the shared MQTT or WebRTC connection.
Package formats and utilities
The npm package provides ESM, CommonJS, TypeScript declarations, and a UMD browser bundle. CDN exports are available through the global Magpie object.
Logger exposes setLevel(), debug(), info(), warning(), and error(). getUtcTimestamp() and getUniqueId() provide the timestamp and identity formats used by frames and RPC.
The TypeScript source is the canonical type reference. See the TypeScript and browser examples for complete programs.