Packages
a2a_elixir_sdk
1.0.1
Elixir implementation of the Agent-to-Agent (A2A) protocol. Exposes ADK agents as A2A-compatible HTTP endpoints and consumes remote A2A agents as local ADK agents.
Current section
Files
Jump to
Current section
Files
lib/a2a_ex/server.ex
defmodule A2AEx.Server do
@moduledoc """
Plug-based HTTP server for the A2A protocol.
Provides endpoints:
- `GET /.well-known/agent-card.json` — Serves the agent card (v0.3.0)
- `GET /.well-known/agent.json` — Serves the agent card (v0.2.5 compat)
- `POST /` — JSON-RPC dispatch (sync and streaming)
## Usage
# Create a handler config
handler = %A2AEx.RequestHandler{
executor: MyExecutor,
task_store: {A2AEx.TaskStore.InMemory, store_pid},
agent_card: %A2AEx.AgentCard{...}
}
# Use as a Plug in your application
plug A2AEx.Server, handler: handler
## With Bandit
Bandit.start_link(plug: {A2AEx.Server, handler: handler}, port: 4000)
[Bandit](https://hex.pm/packages/bandit) is the recommended pure-Elixir
HTTP server. Add `{:bandit, "~> 1.0"}` to your deps.
"""
@behaviour Plug
import Plug.Conn
@impl Plug
def init(opts), do: opts
@impl Plug
def call(%Plug.Conn{} = conn, opts) do
handler = Keyword.fetch!(opts, :handler)
case {conn.method, conn.path_info} do
{"GET", [".well-known", "agent-card.json"]} ->
handle_agent_card(conn, handler)
{"GET", [".well-known", "agent.json"]} ->
handle_agent_card(conn, handler)
{"POST", []} ->
handle_jsonrpc(conn, handler)
{"POST", _} ->
handle_jsonrpc(conn, handler)
_ ->
send_resp(conn, 404, "not found")
end
end
# --- Agent card ---
defp handle_agent_card(conn, handler) do
case handler.agent_card do
nil ->
send_json(conn, 404, %{"error" => "agent card not configured"})
card ->
send_json(conn, 200, A2AEx.AgentCard.to_map(card))
end
end
# --- JSON-RPC dispatch ---
defp handle_jsonrpc(conn, handler) do
{:ok, body, conn} = read_body(conn)
case A2AEx.JSONRPC.decode_request(body) do
{:ok, request} ->
if A2AEx.JSONRPC.streaming_method?(request.method) do
handle_streaming(conn, handler, request)
else
handle_sync(conn, handler, request)
end
{:error, error, req_id} ->
send_jsonrpc_error(conn, error, req_id)
end
end
# --- Sync handling ---
defp handle_sync(conn, handler, request) do
case A2AEx.RequestHandler.handle(handler, request) do
{:ok, result} ->
send_jsonrpc_response(conn, result, request.id)
{:error, error} ->
send_jsonrpc_error(conn, error, request.id)
{:stream, _task_id} ->
error = A2AEx.Error.new(:internal_error, "unexpected stream response")
send_jsonrpc_error(conn, error, request.id)
end
end
# --- Streaming handling (SSE) ---
defp handle_streaming(conn, handler, request) do
case A2AEx.RequestHandler.handle(handler, request) do
{:stream, task_id} ->
conn = start_sse(conn)
stream_events(conn, handler, task_id, request.id)
{:error, error} ->
send_streaming_error(conn, request, error)
{:ok, result} ->
send_jsonrpc_response(conn, result, request.id)
end
end
# tasks/resubscribe errors: send as SSE event per A2A spec
defp send_streaming_error(conn, %{method: "tasks/resubscribe"} = request, error) do
conn = start_sse(conn)
error_resp = A2AEx.JSONRPC.error_map(error, request.id)
send_sse_event(conn, error_resp)
conn
end
# Other streaming method errors (e.g. invalid params): send as JSON
defp send_streaming_error(conn, request, error) do
send_jsonrpc_error(conn, error, request.id)
end
defp start_sse(conn) do
conn
|> put_resp_header("content-type", "text/event-stream")
|> put_resp_header("cache-control", "no-cache")
|> put_resp_header("connection", "keep-alive")
|> send_chunked(200)
end
defp stream_events(conn, handler, task_id, request_id) do
receive do
{:a2a_event, ^task_id, event} ->
task_map = update_and_get_task(handler, task_id, event)
resp = A2AEx.JSONRPC.response_map(task_map, request_id)
case send_sse_event(conn, resp) do
{:ok, conn} -> stream_events(conn, handler, task_id, request_id)
{:error, _} -> conn
end
{:a2a_done, ^task_id} ->
conn
after
60_000 ->
error = A2AEx.Error.new(:internal_error, "stream timeout")
error_resp = A2AEx.JSONRPC.error_map(error, request_id)
send_sse_event(conn, error_resp)
conn
end
end
defp update_and_get_task(handler, task_id, event) do
A2AEx.RequestHandler.update_task_from_event(handler, task_id, event)
case A2AEx.RequestHandler.get_task(handler, task_id) do
{:ok, task} -> A2AEx.Task.to_map(task)
{:error, _} -> A2AEx.Event.to_map(event)
end
end
defp send_sse_event(conn, data) do
json = Jason.encode!(data)
chunk(conn, "data: #{json}\n\n")
end
# --- Response helpers ---
defp send_jsonrpc_response(conn, result, id) do
resp = A2AEx.JSONRPC.response_map(result, id)
send_json(conn, 200, resp)
end
defp send_jsonrpc_error(conn, %A2AEx.Error{} = error, id) do
resp = A2AEx.JSONRPC.error_map(error, id)
send_json(conn, 200, resp)
end
defp send_json(conn, status, data) do
conn
|> put_resp_content_type("application/json")
|> send_resp(status, Jason.encode!(data))
end
end