Interop
Reach agents over A2A, MCP, and AG-UI without changing their messages
This page is the reference for the protocol bridges. For the agents behind them, start with Agents.
Bridges connect external A2A, MCP, and AG-UI clients to agents that speak AGDX. A bridge maps the shared fields into an AGDX command and keeps the protocol's own data in the body. The agents behind a bridge do not change. Bridges work on plain Apache Iggy and need no model provider.
Put a worker behind a bridge
A bridge publishes an AGDX command on its request topic and waits for a correlated AGDX response on its reply topic. The worker is an ordinary agent that listens on the request topic and answers with respond. respond needs the agent's respond_on topic, which must be the bridge's reply topic. It answers an AGDX command with a typed response that carries the correlation and is addressed to the bridge.
import { Agent, AgentId, AgentTopic, agentMessageBody } from "@laserdata/laser-sdk"
const text = new TextEncoder()
const decode = (bytes: Uint8Array) => new TextDecoder().decode(bytes)
await using assistant = Agent.builder()
.id(AgentId.new("assistant"))
.listenOn(AgentTopic.Sessions)
.respondOn(AgentTopic.Sessions)
.handler({
handle: (message, ctx) =>
ctx.respond(text.encode(`handled: ${decode(agentMessageBody(message))}`))
})
.build()
.spawn(laser)
await assistant.ready()use laser_sdk::prelude::full::*;
struct Assistant;
impl AgentHandler for Assistant {
async fn handle(&self, message: &AgentMessage, ctx: &AgentCtx<'_>) -> Result<(), LaserError> {
let text = String::from_utf8_lossy(message.body()).into_owned();
ctx.respond(format!("handled: {text}")).await
}
}
let mut assistant = Agent::builder()
.id("assistant".parse()?)
.listen_on(AgentTopic::Sessions)
.respond_on(AgentTopic::Sessions)
.handler(Assistant)
.build()
.spawn(laser.clone());
assistant.ready().await?;async def assistant(ctx, message):
await ctx.respond(b"handled: " + bytes(message.body()))
worker = laser.spawn_agent(
"assistant",
ls.AgentTopic.Sessions,
assistant,
respond_on=ls.AgentTopic.Sessions,
)
await worker.ready()Bridge constructors take the request topic and the reply topic explicitly. The samples on this page use agent.sessions for both. Every topic lives on the stream of the Laser you pass, so pass laser.with_default_stream(stream) to run a bridge on another stream.
A2A
A2aBridge serves an internal agent through A2A JSON-RPC. card() returns its A2A v1.0 Agent Card.
| A2A method | Mapping |
|---|---|
SendMessage | Publishes an AGDX command on a new task conversation and returns Submitted. The task ID is the conversation ID. The A2A params ride in the body as JSON. |
SendStreamingMessage | Publishes the same command as SendMessage. Readers consume the stream from the log. The bridge does not send it as SSE. |
GetTask | Reads the reply topic and maps the response or error with the task's correlation to the A2A task. Returns Working until a reply arrives. |
CancelTask | Publishes an AGDX error terminal with code Cancelled and returns Canceled. |
The bridge does not serve ListTasks, because it keeps no state outside the log.
import { A2aBridge, AgentId, AgentTopic } from "@laserdata/laser-sdk"
const bridge = new A2aBridge(
laser,
AgentId.new("a2a-gateway"),
AgentTopic.Sessions,
AgentTopic.Sessions
)
const card = bridge.card()use laser_sdk::prelude::full::*;
use std::sync::Arc;
let bridge = Arc::new(A2aBridge::new(
laser.clone(),
"a2a-gateway".parse::<AgentId>()?,
AgentTopic::Sessions,
AgentTopic::Sessions,
));
let card = bridge.card();
let app = bridge.clone().router();import laser_sdk as ls
bridge = ls.A2aBridge(
laser,
"a2a-gateway",
ls.AgentTopic.Sessions,
ls.AgentTopic.Sessions,
)
card = bridge.card()The A2A bridge exposes submit, task, cancel, card, and signed_card (TypeScript: signedCard). handle_rpc(request) (TypeScript: handleRpc) serves a whole JSON-RPC request without HTTP. Rust's a2a-http feature adds router(), an Axum router with the JSON-RPC endpoint at / and the Agent Card at /.well-known/agent-card.json, which delegates to handle_rpc. router() consumes an Arc, so call it on a clone. TypeScript and Python have no router, so route HTTP requests to handle_rpc yourself.
const response = await bridge.handleRpc({
jsonrpc: "2.0",
id: 1,
method: "GetTask",
params: { id: taskId }
})let response = bridge
.handle_rpc(serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "GetTask",
"params": { "id": task_id },
}))
.await;response = await bridge.handle_rpc(
{
"jsonrpc": "2.0",
"id": 1,
"method": "GetTask",
"params": {"id": task_id},
}
)What to expect from the A2A bridge:
submittakes theSendMessageparams as raw JSON bytes in Rust. TypeScript and Python also accept a JSON string, which rides byte for byte, or a value they encode as JSON. Throughhandle_rpcand the router the params are parsed and encoded again, so key order and spacing can change.GetTaskandCancelTaskdo not check that the task exists. An unknown task ID readsWorking.CancelTaskalways publishes and returnsCanceled, but a laterGetTaskreturns the first terminal it finds, so a completed task stays completed.GetTaskscans the reply topic from its start, up to 32 passes of up to 10,000 records per partition each. A reply deeper than that readsWorking.- The task lookup needs a default stream on the
Laser. - A Python task is a dictionary, such as
task["id"]. TypeScript reports the task state as an object, such as{ kind: "known", name: "Completed" }. The JSON-RPC output is the same in all three.
Sign and verify the Agent Card
with_capabilities puts the agent's capability descriptors on the card as A2A skills. signed_card(key) adds a detached JWS signature to the card, so clients can check it before they trust it. with_signing_key signs the CancelTask terminal that the bridge publishes. It signs nothing else, and submitted commands are never signed. TypeScript uses withCapabilities, withSigningKey, and signedCard. Python passes capabilities= and signing_key= to A2aBridge(laser, ..) and calls bridge.signed_card(key).
The Rust router serves the unsigned card. To publish a signed card, serve signed_card yourself.
A client checks a signed card with verify_card. It rebuilds the signed payload from the card itself, so a changed field fails verification.
import { SigningKey, verifyCard } from "@laserdata/laser-sdk"
const key = SigningKey.fromBytes(secret)
const signed = bridge.signedCard(key)
const [signature] = signed.signatures ?? []
if (signature !== undefined) {
verifyCard(signed, signature, key.verifyingKey())
}use laser_sdk::sign::{SigningKey, verify_card};
let key = SigningKey::from_bytes(&secret);
let signed = bridge.signed_card(&key)?;
let value = serde_json::to_value(&signed)
.map_err(|error| LaserError::Codec(error.to_string()))?;
verify_card(&value, &signed.signatures[0], &key.verifying_key())?;key = ls.SigningKey(secret)
signed = bridge.signed_card(key)
ls.verify_card(signed, signed["signatures"][0], key.verifying_key)secret is a 32-byte Ed25519 seed. Rust needs the sign and a2a-bridge features for the card helpers.
MCP
McpBridge maps MCP JSON-RPC tool calls to AGDX commands. It waits for a response or error with the matching correlation and returns it as a tool result.
| MCP method | Mapping |
|---|---|
initialize | Echoes the client's protocol version, 2025-11-25 when the client sends none. Always advertises tools, and advertises resources and prompts only when configured. |
tools/list / tools/call | Serves the tools added with with_tool. A call publishes an AGDX command with the tool name and returns the correlated reply as a tool result. |
resources/list / resources/read | Serves the text resources added with with_resource. |
prompts/list / prompts/get | Serves the prompts added with with_prompt. Arguments are listed but not substituted. |
The bridge serves only these seven methods. Any other method, including notifications such as notifications/initialized and ping, gets an error response. JSON-RPC batch arrays are rejected.
A tool call waits 30 seconds for its reply by default. Change it with with_timeout (TypeScript: withTimeout(ms), Python: timeout_ms=). with_memory_tools adds remember and recall tools with fixed input schemas (TypeScript: withMemoryTools, Python: memory_tools=True).
import { AgentId, AgentTopic, McpBridge } from "@laserdata/laser-sdk"
const mcp = new McpBridge(
laser,
AgentId.new("mcp-gateway"),
AgentTopic.Sessions,
AgentTopic.Sessions,
"ops-tools"
).withTool("ask", "ask the assistant", { type: "object" })
const tools = mcp.listTools()let mcp = Arc::new(
McpBridge::new(
laser.clone(),
"mcp-gateway".parse::<AgentId>()?,
AgentTopic::Sessions,
AgentTopic::Sessions,
"ops-tools",
)
.with_tool(
"ask",
Some("ask the assistant".into()),
serde_json::json!({ "type": "object" }),
)?,
);
let tools = mcp.list_tools();
let app = mcp.clone().router();mcp = ls.McpBridge(
laser,
"mcp-gateway",
ls.AgentTopic.Sessions,
ls.AgentTopic.Sessions,
"ops-tools",
tools=[
{
"name": "ask",
"description": "ask the assistant",
"input_schema": {"type": "object"},
}
],
)
tools = mcp.list_tools()Rust and TypeScript add tools, resources, and prompts with with_tool, with_resource, and with_prompt (TypeScript: withTool, withResource, withPrompt). Python passes tools=, resources=, and prompts= lists of dictionaries. A resource dictionary has uri, name, mime_type, and text. A prompt dictionary has a prompt entry and a messages list of role and text pairs. Every language requires the input schema to be a JSON object. A missing schema, or one that is not an object, is an InvalidError, and Rust with_tool returns a Result for that reason.
Every language also exposes initialize, list_tools, call_tool, list_resources, read_resource, list_prompts, and get_prompt (camelCase in TypeScript). call_tool(name, params) takes the arguments as raw JSON bytes in Rust. Python and TypeScript also take raw JSON text, sent as is, or a value they encode as JSON, such as { question: "status?" }. Through handle_rpc the body is the whole MCP params object, with name and arguments.
A tool that answers with an AGDX error returns a successful JSON-RPC result with isError: true. Each content item names its type in a kind field in Rust and TypeScript, and the JSON-RPC output and Python use type.
Rust's mcp-http feature adds an Axum router(). The mcp-bridge feature alone has no HTTP dependency. handle_rpc(request) (TypeScript: handleRpc) dispatches one JSON-RPC request without HTTP.
Address one agent
A bridge call goes to every agent unless you name one. On the shared agent.sessions topic every listening agent receives an unaddressed task or tool call, so each of them would handle it. Name the agent, and only that agent takes the command as work.
const worker = AgentId.new("assistant")
const task = await bridge.submit(
'{"message":{"role":"user","parts":[{"kind":"text","text":"summarize"}]}}',
{ target: worker }
)
const result = await mcp.callTool("ask", { question: "status?" }, { target: worker })let task = bridge
.submit_to(
"assistant".parse::<AgentId>()?,
br#"{"message":{"role":"user","parts":[{"kind":"text","text":"summarize"}]}}"#.to_vec(),
)
.await?;
let result = mcp
.call_tool_to("assistant".parse::<AgentId>()?, "ask", br#"{"question":"status?"}"#.to_vec())
.await?;task = await bridge.submit(
b'{"message":{"role":"user","parts":[{"kind":"text","text":"summarize"}]}}',
target="assistant",
)
result = await mcp.call_tool("ask", b'{"question":"status?"}', target="assistant")The Rust submit_to and call_tool_to take any agent ID that converts into the wire ID, so the prelude AgentId works directly. Give an inline parse its type, as in "assistant".parse::<AgentId>()?. Requests that arrive through the JSON-RPC routes are unaddressed, so put one worker on the request topic behind them.
Run a bridge call as a child session
A bridge call can run as a child session of a session you hold. The child gets a submitted start on agent.sessions with its parent and root before the command, and the command carries the same ancestry, so a session view rolls the call up under its parent. submit_in leaves the end to the agent that handles the task. If the command send fails after the start, the child stays open. call_tool_in ends the child by the result: completed on a tool result, failed on a tool error, a timeout, or a transport error.
const parent = session.conversation
const root = session.root ?? parent
const task = await bridge.submitIn(
parent,
root,
'{"message":{"role":"user","parts":[{"kind":"text","text":"summarize"}]}}',
{ target: worker }
)
const result = await mcp.callToolIn(parent, root, "ask", { question: "status?" }, { target: worker })let parent = session.conversation();
let root = session.root().unwrap_or(parent);
let task = bridge
.submit_in_to(
"assistant".parse::<AgentId>()?,
parent.into(),
root.into(),
br#"{"message":{"role":"user","parts":[{"kind":"text","text":"summarize"}]}}"#.to_vec(),
)
.await?;
let result = mcp
.call_tool_in_to(
"assistant".parse::<AgentId>()?,
parent.into(),
root.into(),
"ask",
br#"{"question":"status?"}"#.to_vec(),
)
.await?;parent = session.conversation
root = session.root or parent
task = await bridge.submit_in(
parent,
b'{"message":{"role":"user","parts":[{"kind":"text","text":"summarize"}]}}',
root=root,
target="assistant",
)
result = await mcp.call_tool_in(
parent, "ask", b'{"question":"status?"}', root=root, target="assistant"
)Rust also has submit_in and call_tool_in without a target. Python and TypeScript make the target optional. Child lifecycle records carry neither the loop guard nor a signature.
AG-UI
AG-UI gives frontends state and events. publish_state_snapshot and publish_state_delta write the session state document of a conversation, the same records that session.state() writes, so every reader folds one state model. A snapshot must be a JSON object, and a delta must be an RFC 6902 JSON Patch array. reconstruct_state returns that document. It folds agent.sessions first, and reads the managed state view only when the lane no longer holds the document's baseline and the view is not behind the lane. It returns nothing until a state record exists.
agui_events reads a conversation from the log and converts it into AG-UI events:
- Chat chunks become
TEXT_MESSAGE_*events, and reasoning chunks becomeREASONING_MESSAGE_*events. - Tool argument chunks become
TOOL_CALL_START,TOOL_CALL_ARGS, andTOOL_CALL_END. A response or error that names a tool becomesTOOL_CALL_RESULT. - A task status of
SubmittedbecomesRUN_STARTED. Completed or canceled status becomesRUN_FINISHED. Failed or rejected status becomesRUN_ERROR, with the status detail ortask failed. Another error also becomesRUN_ERROR. Only status records with operationtaskrender. - State records become
STATE_SNAPSHOTandSTATE_DELTA.
AG-UI events with no AGDX source, such as MESSAGES_SNAPSHOT, ACTIVITY_*, RAW, and CUSTOM, are never produced.
const conversation = ConversationId.new()
await laser.publishStateSnapshot(AgentId.new("ui"), conversation, { count: 0 })
const events = await laser.aguiEvents(conversation, AgentTopic.Sessions)let conversation = ConversationId::new();
laser
.publish_state_snapshot("ui".parse::<AgentId>()?, conversation, &serde_json::json!({ "count": 0 }))
.await?;
let events = laser.agui_events(conversation, AgentTopic::Sessions).await?;conversation_id = ls.new_conversation_id()
await laser.publish_state_snapshot("ui", conversation_id, {"count": 0})
events = await laser.agui_events(conversation_id, ls.AgentTopic.Sessions)Pause for a human decision
request_input publishes a prompt command and waits for the correlated response on the reply topic, up to the timeout. The approver answers with respond_input. If the approver rejects with an error, the wait fails with LaserError::Rejected in Rust and RejectedError in TypeScript and Python. Both calls use the existing command and response records and need no bridge.
Name the approver so that only it receives the prompt. Without a target every agent on agent.sessions receives it, and any of them could answer. Rust uses request_input_from(target, ..), Python target=, and TypeScript a { target } option.
const decision = await laser
.agdx(AgentTopic.Sessions, AgentId.new("orchestrator"), conversation)
.requestInput(
AgentTopic.Sessions,
new TextEncoder().encode("restart the metrics service on node-7?"),
300_000,
{ target: AgentId.new("approver") }
)let decision = laser
.agdx(AgentTopic::Sessions, "orchestrator".parse::<AgentId>()?, conversation.into())
.request_input_from(
"approver".parse::<AgentId>()?,
AgentTopic::Sessions,
b"restart the metrics service on node-7?".to_vec(),
Duration::from_secs(300),
)
.await?;decision = await laser.agdx(
ls.AgentTopic.Sessions,
"orchestrator",
conversation_id,
).request_input(
ls.AgentTopic.Sessions,
b"restart the metrics service on node-7?",
timeout_ms=300_000,
target="approver",
)The timeout is required in every SDK and has no default. Python takes it as timeout_ms= in milliseconds. TypeScript also accepts a signal option, and an aborted signal throws CancelledError.
The approver is an agent that listens on agent.sessions. Its handler answers with respond_input, which needs the handled message to be an AGDX command with a correlation and the agent to have an ID. An agent with a signing key signs the response, so a caller that verifies signatures accepts only that agent's decision. request_input does not sign the prompt.
const approver = {
handle: (_message: AgentMessage, ctx: AgentCtx) =>
ctx.respondInput(AgentTopic.Sessions, new TextEncoder().encode("approved"))
}struct Approver;
impl AgentHandler for Approver {
async fn handle(&self, _message: &AgentMessage, ctx: &AgentCtx<'_>) -> Result<(), LaserError> {
ctx.respond_input(AgentTopic::Sessions, b"approved".to_vec()).await
}
}async def approver(ctx, message):
await ctx.respond_input(ls.AgentTopic.Sessions, b"approved")Inside a handler, ctx.approval_gate(reply_topic, prompt, timeout) (TypeScript: approvalGate) does the same from the handled conversation and returns the decision body. It always publishes the prompt unaddressed on agent.sessions and has no target option.
Authorize edge requests
The JSON-RPC bridges do not authenticate HTTP requests. Add authentication in your HTTP framework. Your HTTP layer decodes the token, and authorize_edge (TypeScript: authorizeEdge) checks the audience and scope of the decoded claims. The SDK does not parse tokens.
import { authorizeEdge, edgeDenialChallenge } from "@laserdata/laser-sdk"
const denial = authorizeEdge(
{ audience: ["agents.example"], scopes: ["tool:read"] },
"agents.example",
"tool:write"
)
if (denial !== undefined) {
const challenge = edgeDenialChallenge(denial)
}use laser_sdk::edge_auth::{EdgeClaims, authorize_edge};
let claims = EdgeClaims {
audience: vec!["agents.example".to_owned()],
scopes: vec!["tool:read".to_owned()],
};
if let Err(denial) = authorize_edge(&claims, "agents.example", "tool:write") {
let challenge = denial.challenge();
}denial = ls.authorize_edge(
["agents.example"],
["tool:read"],
"agents.example",
"tool:write",
)
if denial is not None:
challenge = denial.challengeA token for another audience is a hard reject with no challenge, which maps to 401. A token for the right audience without the required scope returns the challenge Bearer scope="tool:write", which maps to 403 with a step-up. Rust returns the denial as the error, with expected or required_scope on its variant. TypeScript returns the denial or undefined, with kind wrongAudience or stepUp and the fields expected and requiredScope. Python returns an EdgeDenial with kind (wrong_audience or step_up), code, challenge, expected, and required_scope, or None. In Rust, authorize_edge needs the a2a-bridge or mcp-bridge feature. Governance covers delegated grants.
with_default_stream (TypeScript: withDefaultStream) points a bridge at another stream on the same connection. Use separate credentials when Iggy permissions must keep the streams apart.
Loop guard
Each bridge stamps its own ID in the bridge_hops metadata of the commands and cancels it publishes. Use with_bridge_hops(previous) (TypeScript: withBridgeHops) to continue a path that a call already took through other bridges. A path that already holds the bridge's ID is rejected before anything is published. The guard never reads inbound records, so pass previous yourself. Rust's with_bridge_hops consumes the bridge and returns a result. TypeScript and Python change the bridge in place, and Python refuses while a call is in flight.
Errors
A failed bridge call still returns a JSON-RPC response envelope. Its error carries the application code -32000 and a public message from a fixed set: invalid request, unsupported operation, conflict, unauthenticated, forbidden, step-up authorization required, not found, or internal error. A malformed request gets invalid request. A tool call timeout reads internal error. Over the Rust routers, a body that is not JSON is rejected with an HTTP 4xx before handle_rpc runs.
With a key registry verifier on the Laser (Rust sign feature), GetTask, MCP calls, and request_input skip replies that fail verification.
Standalone helpers
These helpers need no connection and behave the same in every language.
| Helper | Rust | TypeScript | Python |
|---|---|---|---|
| A2A request to AGDX command | a2a::command_from_message_send | commandFromMessageSend | command_from_message_send |
| AGDX envelope to A2A task | a2a::task_from_envelope | taskFromEnvelope | task_from_envelope |
| MCP tool call to AGDX command | mcp::tool_call_from_request | toolCallFromRequest | tool_call_from_request |
| AGDX envelope to MCP tool result | mcp::tool_result_from_envelope | toolResultFromEnvelope | tool_result_from_envelope |
| Bridge loop guard | a2a::enter_bridge | enterBridge | enter_bridge |
| Edge token check | edge_auth::authorize_edge | authorizeEdge | authorize_edge |
| Sign an Agent Card | sign::sign_card_value | signCardValue | sign_card_value |
| Verify an Agent Card | sign::verify_card | verifyCard | verify_card |
| Verify a signed delegation | sign::verify_delegation | verifyDelegation | verify_delegation |
| Claim check in and out | blob::check_in, blob::resolve_body | checkIn, resolveBody | check_in, resolve_body |
| Snapshot bytes | snapshot::encode, snapshot::decode | encodeSnapshot, decodeSnapshot | encode_snapshot, decode_snapshot |
| Resume offsets of a snapshot | agent::resume_offsets | resumeOffsets | resume_offsets |
| Reciprocal-rank fusion | memory::fuse_reciprocal_rank | fuseReciprocalRank | fuse_reciprocal_rank |
| Clocks | agent::SystemClock, agent::TestClock | SystemClock, TestClock | SystemClock, TestClock |
- The converters keep the original JSON bytes. Set the JSON content type when you publish the resulting envelope.
enter_bridgerejects an empty bridge ID and a path that already contains this bridge. In Rust it needs thea2a-bridgeor themcp-bridgefeature.verify_delegationchecks that the envelope signature verifies against an enrolled key and returns the signer and the delegated user from the signed metadata, or nothing when the envelope carries no delegation. Authorize the pair withdelegated_allow.resume_offsetsreturns one past the last offset a snapshot folded for each partition, and stops at the unsigned 64-bit ceiling. The snapshot byte helpers read and write the named-field CBOR storage form, and fail with an invalid error on a bad snapshot.SystemClockandTestClockcount epoch microseconds as unsigned 64-bit values, andadvancewraps around at the ceiling. A TypeScriptTestClockneeds abigintstart, and Python defaults the start to 0. Both raiseInvalidErroroutside the range.- The card helpers need the
signanda2a-bridgefeatures in Rust. Thesignfeature alone coversverify_delegation.
Where it runs
Bridges work on plain Apache Iggy, Laser Stack, and LaserData Cloud. reconstruct_state reads the managed state view only on a deployment that has one. The SDK does not host a model loop or a sandbox. A provider's event subscription is not an Iggy stream, so an adapter in front of one must map event IDs and cursors, record tool results with their correlation, and collect child results explicitly.
In Rust, enable a2a-bridge, mcp-bridge, or agui for the adapters, and add a2a-http or mcp-http for the Axum routers.