A2aEngine

Pure-spec implementation of the Agent2Agent (A2A) protocol — types, JSON-RPC 2.0 codec, SSE codec, and a transport behaviour.

A2A defines how agents call each other: JSON-RPC 2.0 envelopes over HTTPS (with SSE for streaming), a typed Task / Message / Part / Artifact model, AgentCard discovery, and push-notification config. A2aEngine implements that wire contract in Elixir and nothing more.

It is deliberately stateless. There is no database, no broker, no process registry, no opinion about where a Task lives or how it is resumed. The host application owns all state — this library gives it the vocabulary and the plumbing to move A2A messages over any transport it chooses.

Why this shape?

Most protocol libraries bundle a server, a store, and a client. That's convenient until two of your agents already live in the same BEAM cluster and you'd rather not pay the HTTP/JSON round-trip. A2aEngine splits the contract from the substrate:

Bring your own handler, your own auth, your own task store — wire them into a transport and you speak A2A.

Installation

Add a2a_engine to your dependencies:

def deps do
[
{:a2a_engine, "~> 0.1.0"}
]
end

Documentation is published at https://hexdocs.pm/a2a_engine.

What's in the box

AreaModules
Spec typesA2aEngine.Types.*Task, Message, Part (Text/File/Data), Artifact, AgentCard, PushNotificationConfig, lifecycle events, …
CodecCodec.JsonRpc (envelopes + A2A method registry), Codec.SSE (encode/decode + Last-Event-ID resumption), Codec.Keys (camelCase ↔ snake_case)
TransportTransport (behaviour), Transport.Http (+ .Plug), Transport.BeamNative (+ .Server)
AuthAuth (behaviour), Auth.Bearer (shared-secret), Auth.Localhost (same-host bypass)
CoreHandler (request/stream callback behaviour), Errors (A2A JSON-RPC codes), LifecycleState (task state enum), Parts (checkpoint-owner filtering)
TestingTestPeer — a scriptable in-process A2A peer with a real HTTP server, shipped in lib/ so hosts can reuse it

Quick tour

JSON-RPC envelopes

alias A2aEngine.Codec.JsonRpc
# Build an outbound request
request =
JsonRpc.encode_request("message/send", %{"message" => msg}, "req-1")
# => %{"jsonrpc" => "2.0", "id" => "req-1",
# "method" => "message/send", "params" => %{"message" => msg}}
# Decode an inbound request
{:ok, %{id: id, method: method, params: params}} =
JsonRpc.decode_request(envelope)
# Decode a response (success or error)
{:ok, {:success, id, result}} = JsonRpc.decode_response(envelope)
{:ok, {:error, id, %{code: code, message: msg}}} = JsonRpc.decode_response(envelope)
# The A2A method registry
:true = JsonRpc.known_method?("message/send")
:message_send = JsonRpc.atom_for_method("message/send")

params and result pass through as plain maps — encode/decode the A2A types (via A2aEngine.Types.*) on either side of the envelope.

SSE streaming with resumption

alias A2aEngine.Codec.SSE
# Encode an event frame (id scheme: "<task_id>:<seq>")
frame = SSE.encode_event("task-42", "taskstatus", 3, %{state: "working"})
# => "id: task-42:3\nevent: taskstatus\ndata: {...}\n\n"
# Resume a stream from the Last-Event-ID header
{:ok, {"task-42", 3}} = SSE.resume_from(last_event_id)
{:error, :no_last_event_id} = SSE.resume_from(nil)

Implement a handler

A handler is a behaviour that fulfils JSON-RPC methods. It is transport-agnostic — both transports dispatch to the same callbacks.

defmodule MyApp.A2A.Broker do
@behaviour A2aEngine.Handler
alias A2aEngine.Errors
@impl true
def handle_request("tasks/get", %{"id" => id}, _ctx) do
{:ok, fetch_task!(id)}
end
def handle_request("message/send", params, ctx) do
# ctx.auth carries the authenticated principal
{:ok, run_message(params, ctx.auth)}
end
def handle_request(_method, _params, _ctx) do
{:error, Errors.build(Errors.method_not_found())}
end
@impl true
def handle_stream("message/stream", params, ctx, emit) do
emit.(%{kind: "taskstatus", status: %{state: "working"}})
emit.(%{kind: "taskstatus", status: %{state: "completed"}})
{:ok, :done}
end
end

Call an agent over HTTP

alias A2aEngine.{Codec.JsonRpc, Transport.Http}
request = JsonRpc.encode_request("message/send", %{"message" => msg}, "req-1")
# Unary
{:ok, response} =
Http.send_request("https://agent.example/a2a", request,
headers: [{"authorization", "Bearer " <> token}]
)
# Streaming — returns a lazy Enumerable of decoded events
{:ok, stream} = Http.stream_request("https://agent.example/a2a", request)
for event <- stream, do: handle_event(event)

Serve A2A from a Plug

# in your router or Endpoint
forward "/a2a", A2aEngine.Transport.Http.Plug,
handler: MyApp.A2A.Broker,
auth: A2aEngine.Auth.Bearer,
auth_opts: [tokens: %{"secret" => "my-agent"}]

The host must parse the JSON body upstream (e.g. via Plug.Parsers). The plug reads conn.body_params, dispatches to the handler, and returns a JSON-RPC response — or an SSE stream for message/stream / tasks/resubscribe.

Test against a scriptable peer

TestPeer spins up a real Bandit HTTP server on a random port and lets you script responses per method. Ship your integration tests without standing up a fake agent by hand.

{:ok, peer} = A2aEngine.TestPeer.start_link()
url = A2aEngine.TestPeer.url(peer)
A2aEngine.TestPeer.script(peer, "message/send", fn _params ->
{:ok, %{"kind" => "task", "id" => "t1", "status" => %{"state" => "completed"}}}
end)
{:ok, response} = A2aEngine.Transport.Http.send_request(url, request)

Design notes

License

MIT. See LICENSE.