Add a2a module: agent card, tasks, streaming, input-required and agent-to-agent calls on the A2A Java SDK, served over the JDK HTTP server

Co-Authored-By: Claude Sonnet 5.5 <[email protected]>
Claude-Session: https://claude.ai/code/session_01G8ikz8xdWuTP5yun8DZ1hk
This commit is contained in:
Claude
2026-10-11 08:35:19 +00:00
parent ec1b076e9b
commit 10545299f8
19 changed files with 1042 additions and 2 deletions
+4 -2
View File
@@ -7,12 +7,14 @@ facts, so the build fails when a claim stops being true.
| Module | Post | What it shows |
|---|---|---|
| `embabel` | Embabel: Goal-Oriented AI Agents on the JVM | actions, goals, conditions and cost-based planning with GOAP, compared with plain Spring AI |
| `a2a` | A2A Protocol in Java: Agents That Talk to Each Other | agent card, tasks, streaming, input-required and agent-to-agent calls with the A2A Java SDK |
## Versions (verified on Maven Central, 2026-10-11)
| Component | Version | Note |
|---|---|---|
| Embabel `embabel-agent-api` | 1.5.3 | brings Spring AI 2.0.1 and Spring Boot 4.1.1 transitively |
| A2A Java SDK `io.github.a2asdk` | 1.0.0.Alpha3 | alpha; targets A2A protocol 1.0. The newest stable release, 0.3.3.Final, targets the 0.3 protocol and its API differs |
| JDK | 25 LTS | |
| JUnit | 6.1.3 | |
@@ -21,9 +23,9 @@ facts, so the build fails when a claim stops being true.
```bash
export JAVA_HOME=/path/to/jdk-25
mvn test # everything
mvn test -pl embabel -am # one module
mvn test -pl a2a -am # one module
```
No API key is needed. The Embabel tests use `ScriptedLlmOperations`, the scripted stand-in shipped inside
No API key is needed. The A2A agents contain no model call, because the protocol is the subject. The Embabel tests use `ScriptedLlmOperations`, the scripted stand-in shipped inside
`embabel-agent-api`, and the Spring AI comparison uses a small scripted `ChatModel`.
Each module writes its transcripts to `<module>/output/` when you run the tests.
+40
View File
@@ -0,0 +1,40 @@
GET /.well-known/agent-card.json
{
"name": "Glossary Agent",
"description": "Defines Java terms",
"version": "1.0.0",
"capabilities": {
"streaming": true,
"pushNotifications": false,
"extendedAgentCard": false
},
"defaultInputModes": [
"text/plain"
],
"defaultOutputModes": [
"text/plain"
],
"skills": [
{
"id": "define",
"name": "Define a term",
"description": "Returns a one-sentence definition",
"tags": [
"java",
"glossary"
],
"examples": [
"record"
]
}
],
"supportedInterfaces": [
{
"protocolBinding": "JSONRPC",
"url": "http://127.0.0.1:PORT",
"tenant": "",
"protocolVersion": "1.0"
}
]
}
+61
View File
@@ -0,0 +1,61 @@
POST /
{
"jsonrpc": "2.0",
"id": "1",
"method": "SendMessage",
"params": {
"message": {
"messageId": "m-1",
"role": "ROLE_USER",
"parts": [
{
"text": "record"
}
]
},
"configuration": {
"blocking": true
}
}
}
Response:
{
"jsonrpc": "2.0",
"id": 1,
"result": {
"task": {
"id": "<id-1>",
"contextId": "<id-2>",
"status": {
"state": "TASK_STATE_COMPLETED",
"timestamp": "<time>"
},
"artifacts": [
{
"artifactId": "definition",
"name": "definition",
"description": "",
"parts": [
{
"text": "A record is a final class that is a transparent",
"metadata": {},
"filename": "",
"mediaType": ""
},
{
"text": " carrier for its components.",
"metadata": {},
"filename": "",
"mediaType": ""
}
],
"metadata": {},
"extensions": []
}
],
"history": [],
"metadata": {}
}
}
}
+15
View File
@@ -0,0 +1,15 @@
Same question, "record", asked twice.
Blocking (streaming off): 1 event
TaskEvent state=COMPLETED artifacts=1
answer: A record is a final class that is a transparent carrier for its components.
Streaming: 5 events
TaskStatusUpdate state=SUBMITTED final=false
TaskStatusUpdate state=WORKING final=false
TaskArtifactUpdate append=false lastChunk=false text="A record is a final class that is a transparent"
TaskArtifactUpdate append=true lastChunk=true text=" carrier for its components."
TaskStatusUpdate state=COMPLETED final=true
answer: A record is a final class that is a transparent carrier for its components.
GetTask for the streamed task afterwards: state TASK_STATE_COMPLETED, artifacts 1, artifact text "A record is a final class that is a transparent carrier for its components."
+9
View File
@@ -0,0 +1,9 @@
Turn 1: message with no term.
task <id-1> state INPUT_REQUIRED
agent says: Which term should I define?
Turn 2: message "sealed class" sent with taskId <id-1>.
task <id-1> state COMPLETED
answer: A sealed class restricts which other classes may extend it.
Same task id both turns: true
+17
View File
@@ -0,0 +1,17 @@
Caller -> Brief Agent -> Glossary Agent.
Request: "record, sealed class"
Brief Agent task state: COMPLETED
result:
- record: A record is a final class that is a transparent carrier for its components.
- sealed class: A sealed class restricts which other classes may extend it.
Calls the Brief Agent made to the Glossary Agent:
SendMessage "record" -> COMPLETED "A record is a final class that is a transparent carrier for its components."
SendMessage "sealed class" -> COMPLETED "A sealed class restricts which other classes may extend it."
Request: "record, monad"
Brief Agent task state: FAILED
message: Glossary Agent could not help: I do not know the term 'monad'.
Calls the Brief Agent made to the Glossary Agent:
SendMessage "record" -> COMPLETED "A record is a final class that is a transparent carrier for its components."
SendMessage "monad" -> FAILED "I do not know the term 'monad'."
+42
View File
@@ -0,0 +1,42 @@
mvn dependency:tree for the a2a module (compile scope):
com.ankurm.agents:a2a:jar:1.0.0
io.github.a2asdk:a2a-java-sdk-server-common:jar:1.0.0.Alpha3:compile
io.github.a2asdk:a2a-java-sdk-spec:jar:1.0.0.Alpha3:compile
io.github.a2asdk:a2a-java-sdk-jsonrpc-common:jar:1.0.0.Alpha3:compile
io.github.a2asdk:a2a-java-sdk-spec-grpc:jar:1.0.0.Alpha3:compile
com.google.protobuf:protobuf-java:jar:4.33.1:compile
com.google.api.grpc:proto-google-common-protos:jar:2.63.1:compile
io.github.a2asdk:a2a-java-sdk-http-client:jar:1.0.0.Alpha3:compile
io.smallrye.reactive:mutiny-zero:jar:1.1.1:compile
jakarta.enterprise:jakarta.enterprise.cdi-api:jar:4.1.0:compile
jakarta.enterprise:jakarta.enterprise.lang-model:jar:4.1.0:compile
jakarta.annotation:jakarta.annotation-api:jar:3.0.0:compile
jakarta.el:jakarta.el-api:jar:6.0.0:compile
jakarta.interceptor:jakarta.interceptor-api:jar:2.2.0:compile
jakarta.inject:jakarta.inject-api:jar:2.0.1:compile
org.slf4j:slf4j-api:jar:2.0.17:compile
io.github.a2asdk:a2a-java-sdk-transport-jsonrpc:jar:1.0.0.Alpha3:compile
com.google.protobuf:protobuf-java-util:jar:4.32.1:compile
com.google.code.findbugs:jsr305:jar:3.0.2:runtime
com.google.errorprone:error_prone_annotations:jar:2.18.0:compile
com.google.guava:guava:jar:32.0.1-jre:runtime
com.google.guava:failureaccess:jar:1.0.1:runtime
com.google.guava:listenablefuture:jar:9999.0-empty-to-avoid-conflict-with-guava:runtime
com.google.j2objc:j2objc-annotations:jar:2.8:runtime
com.google.code.gson:gson:jar:2.13.2:compile
io.github.a2asdk:a2a-java-sdk-client:jar:1.0.0.Alpha3:compile
io.github.a2asdk:a2a-java-sdk-client-transport-spi:jar:1.0.0.Alpha3:compile
io.github.a2asdk:a2a-java-sdk-common:jar:1.0.0.Alpha3:compile
io.github.a2asdk:a2a-java-sdk-client-transport-jsonrpc:jar:1.0.0.Alpha3:compile
org.slf4j:slf4j-simple:jar:2.0.17:test
org.junit.jupiter:junit-jupiter:jar:6.1.3:test
org.junit.jupiter:junit-jupiter-api:jar:6.1.3:test
org.opentest4j:opentest4j:jar:1.3.0:test
org.junit.platform:junit-platform-commons:jar:6.1.3:test
org.apiguardian:apiguardian-api:jar:1.1.2:test
org.jspecify:jspecify:jar:1.0.0:test
org.junit.jupiter:junit-jupiter-params:jar:6.1.3:test
org.junit.jupiter:junit-jupiter-engine:jar:6.1.3:test
org.junit.platform:junit-platform-engine:jar:6.1.3:test
org.assertj:assertj-core:jar:3.27.6:test
net.bytebuddy:byte-buddy:jar:1.17.7:test
+5
View File
@@ -0,0 +1,5 @@
The agent accepts the task at once, then works for 400 ms before finishing.
SendMessage with configuration.blocking = false: task state TASK_STATE_WORKING, returned in under 400 ms: true
SendMessage with configuration.blocking = true: task state TASK_STATE_COMPLETED, returned in at least 400 ms: true
SendMessage with no configuration at all: task state TASK_STATE_WORKING, returned in under 400 ms: true
+54
View File
@@ -0,0 +1,54 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>com.ankurm.agents</groupId>
<artifactId>java-ai-agents</artifactId>
<version>1.0.0</version>
</parent>
<artifactId>a2a</artifactId>
<dependencies>
<dependency>
<groupId>io.github.a2asdk</groupId>
<artifactId>a2a-java-sdk-server-common</artifactId>
<version>${a2a.sdk.version}</version>
</dependency>
<dependency>
<groupId>io.github.a2asdk</groupId>
<artifactId>a2a-java-sdk-transport-jsonrpc</artifactId>
<version>${a2a.sdk.version}</version>
</dependency>
<dependency>
<groupId>io.github.a2asdk</groupId>
<artifactId>a2a-java-sdk-client</artifactId>
<version>${a2a.sdk.version}</version>
</dependency>
<dependency>
<groupId>io.github.a2asdk</groupId>
<artifactId>a2a-java-sdk-client-transport-jsonrpc</artifactId>
<version>${a2a.sdk.version}</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId>
<version>${slf4j.version}</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.assertj</groupId>
<artifactId>assertj-core</artifactId>
<version>${assertj.version}</version>
<scope>test</scope>
</dependency>
</dependencies>
</project>
@@ -0,0 +1,160 @@
package com.ankurm.agents.a2a;
import io.a2a.A2A;
import io.a2a.client.Client;
import io.a2a.client.ClientEvent;
import io.a2a.client.MessageEvent;
import io.a2a.client.TaskEvent;
import io.a2a.client.TaskUpdateEvent;
import io.a2a.client.config.ClientConfig;
import io.a2a.client.transport.jsonrpc.JSONRPCTransport;
import io.a2a.client.transport.jsonrpc.JSONRPCTransportConfigBuilder;
import io.a2a.spec.AgentCard;
import io.a2a.spec.Artifact;
import io.a2a.spec.Message;
import io.a2a.spec.Part;
import io.a2a.spec.Task;
import io.a2a.spec.TaskArtifactUpdateEvent;
import io.a2a.spec.TaskState;
import io.a2a.spec.TaskStatusUpdateEvent;
import io.a2a.spec.TextPart;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
/** One call to a remote A2A agent: fetch its card, send a message, collect what comes back until the task settles. */
public final class A2aCaller {
/** What came back: every event in order, and the last task snapshot seen. */
public record Reply(List<ClientEvent> events, Task task, Message message) {
public String text() {
if (task != null && task.artifacts() != null && !task.artifacts().isEmpty()) {
StringBuilder sb = new StringBuilder();
for (Artifact a : task.artifacts()) {
for (Part<?> p : a.parts()) {
if (p instanceof TextPart t) sb.append(t.text());
}
}
return sb.toString();
}
if (task != null && task.status().message() != null) {
return textOf(task.status().message());
}
return message == null ? "" : textOf(message);
}
}
private A2aCaller() {
}
public static AgentCard fetchCard(String agentUrl) throws Exception {
return A2A.getAgentCard(agentUrl);
}
public static Reply send(String agentUrl, String text, String continueTaskId, boolean streaming) throws Exception {
AgentCard card = fetchCard(agentUrl);
List<ClientEvent> events = Collections.synchronizedList(new ArrayList<>());
AtomicReference<Task> task = new AtomicReference<>();
AtomicReference<Message> message = new AtomicReference<>();
AtomicReference<Throwable> error = new AtomicReference<>();
CountDownLatch settled = new CountDownLatch(1);
// The client cancels the SSE request itself once it sees a final state; that cancel comes back as an
// error, so an error only counts when no final state was seen first.
java.util.concurrent.atomic.AtomicBoolean seenFinal = new java.util.concurrent.atomic.AtomicBoolean();
Client client = Client.builder(card)
.withTransport(JSONRPCTransport.class, new JSONRPCTransportConfigBuilder())
.clientConfig(ClientConfig.builder().setStreaming(streaming).build())
.addConsumer((event, c) -> {
events.add(event);
if (event instanceof TaskEvent te) {
task.set(te.getTask());
if (stops(te.getTask().status().state())) { seenFinal.set(true); settled.countDown(); }
} else if (event instanceof TaskUpdateEvent tu) {
task.set(tu.getTask());
if (stops(tu.getTask().status().state())) { seenFinal.set(true); settled.countDown(); }
} else if (event instanceof MessageEvent me) {
message.set(me.getMessage());
seenFinal.set(true);
settled.countDown();
}
})
.streamingErrorHandler(t -> {
error.set(t);
settled.countDown();
})
.build();
try {
Message.Builder m = Message.builder().role(Message.Role.ROLE_USER)
.parts(new TextPart(text)).messageId(UUID.randomUUID().toString());
if (continueTaskId != null) {
m.taskId(continueTaskId);
}
client.sendMessage(m.build());
if (!settled.await(20, TimeUnit.SECONDS)) {
throw new IllegalStateException("Timed out waiting for " + agentUrl);
}
if (error.get() != null && !seenFinal.get()) {
throw new IllegalStateException("Streaming failed", error.get());
}
return new Reply(List.copyOf(events), task.get(), message.get());
} finally {
client.close();
}
}
/** Asks the remote agent for a task it already knows about. */
public static Task getTask(String agentUrl, String taskId) throws Exception {
Client client = Client.builder(fetchCard(agentUrl))
.withTransport(JSONRPCTransport.class, new JSONRPCTransportConfigBuilder())
.build();
try {
return client.getTask(new io.a2a.spec.TaskQueryParams(taskId), null);
} finally {
client.close();
}
}
private static boolean stops(TaskState s) {
return s.isFinal() || s == TaskState.TASK_STATE_INPUT_REQUIRED || s == TaskState.TASK_STATE_AUTH_REQUIRED;
}
static String textOf(Message m) {
StringBuilder sb = new StringBuilder();
for (Part<?> p : m.parts()) {
if (p instanceof TextPart t) sb.append(t.text());
}
return sb.toString();
}
/** A one-line description of an event, for transcripts. */
public static String describe(ClientEvent e) {
if (e instanceof TaskEvent te) {
return "TaskEvent state=" + short_(te.getTask().status().state()) + " artifacts=" + te.getTask().artifacts().size();
}
if (e instanceof MessageEvent me) {
return "MessageEvent text=\"" + textOf(me.getMessage()) + "\"";
}
if (e instanceof TaskUpdateEvent tu) {
if (tu.getUpdateEvent() instanceof TaskStatusUpdateEvent s) {
return "TaskStatusUpdate state=" + short_(s.status().state()) + " final=" + s.isFinal();
}
if (tu.getUpdateEvent() instanceof TaskArtifactUpdateEvent a) {
StringBuilder sb = new StringBuilder();
for (Part<?> p : a.artifact().parts()) if (p instanceof TextPart t) sb.append(t.text());
return "TaskArtifactUpdate append=" + a.append() + " lastChunk=" + a.lastChunk() + " text=\"" + sb + "\"";
}
}
return e.getClass().getSimpleName();
}
static String short_(TaskState s) {
return s.name().replace("TASK_STATE_", "");
}
}
@@ -0,0 +1,224 @@
package com.ankurm.agents.a2a;
import com.sun.net.httpserver.HttpExchange;
import com.sun.net.httpserver.HttpServer;
import io.a2a.grpc.utils.JSONRPCUtils;
import io.a2a.grpc.utils.ProtoUtils;
import io.a2a.jsonrpc.common.json.JsonUtil;
import io.a2a.jsonrpc.common.wrappers.*;
import io.a2a.server.ServerCallContext;
import io.a2a.server.agentexecution.AgentExecutor;
import io.a2a.server.auth.UnauthenticatedUser;
import io.a2a.server.events.InMemoryQueueManager;
import io.a2a.server.events.MainEventBus;
import io.a2a.server.events.MainEventBusProcessor;
import io.a2a.server.requesthandlers.DefaultRequestHandler;
import io.a2a.server.tasks.BasePushNotificationSender;
import io.a2a.server.tasks.InMemoryPushNotificationConfigStore;
import io.a2a.server.tasks.InMemoryTaskStore;
import io.a2a.server.util.sse.SseFormatter;
import io.a2a.spec.A2AError;
import io.a2a.spec.AgentCard;
import io.a2a.spec.InternalError;
import io.a2a.spec.TransportProtocol;
import io.a2a.transport.jsonrpc.handler.JSONRPCHandler;
import java.io.IOException;
import java.io.OutputStream;
import java.net.InetSocketAddress;
import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Flow;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Function;
/**
* Serves one A2A agent over JSON-RPC using only the JDK's built-in HTTP server. The SDK ships the protocol
* logic (tasks, queues, streaming) and a Quarkus binding; this class is the small binding that replaces Quarkus.
* The request routing follows the SDK's own Quarkus route class.
*/
public final class A2aHttpServer implements AutoCloseable {
private final HttpServer http;
private final JSONRPCHandler rpc;
private final AgentCard card;
private final ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor();
private final MainEventBusProcessor processor;
private A2aHttpServer(HttpServer http, JSONRPCHandler rpc, AgentCard card, MainEventBusProcessor processor) {
this.http = http;
this.processor = processor;
this.rpc = rpc;
this.card = card;
}
/** Starts on a free port. The card is built after the port is known, so it can carry the real URL. */
public static A2aHttpServer start(Function<String, AgentCard> cardForUrl, AgentExecutor agent) throws IOException {
HttpServer http = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
String url = "http://127.0.0.1:" + http.getAddress().getPort();
AgentCard card = cardForUrl.apply(url);
ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor();
MainEventBus bus = new MainEventBus();
InMemoryTaskStore taskStore = new InMemoryTaskStore();
InMemoryQueueManager queues = new InMemoryQueueManager(taskStore, bus);
InMemoryPushNotificationConfigStore pushConfigs = new InMemoryPushNotificationConfigStore();
MainEventBusProcessor processor = new MainEventBusProcessor(bus, taskStore, new BasePushNotificationSender(pushConfigs), queues);
callPackagePrivate(processor, "start");
DefaultRequestHandler handler = DefaultRequestHandler.create(agent, taskStore, queues, pushConfigs, processor, executor, executor);
JSONRPCHandler rpc = new JSONRPCHandler(card, handler, executor);
A2aHttpServer server = new A2aHttpServer(http, rpc, card, processor);
http.createContext("/", server::route);
http.setExecutor(server.executor);
http.start();
return server;
}
public String url() {
return "http://127.0.0.1:" + http.getAddress().getPort();
}
public AgentCard card() {
return card;
}
private void route(HttpExchange x) throws IOException {
try (x) {
String path = x.getRequestURI().getPath();
if (x.getRequestMethod().equals("GET") && path.equals("/.well-known/agent-card.json")) {
reply(x, JsonUtil.toJson(card));
} else if (x.getRequestMethod().equals("POST") && path.equals("/")) {
jsonRpc(x, new String(x.getRequestBody().readAllBytes(), StandardCharsets.UTF_8));
} else {
x.sendResponseHeaders(404, -1);
}
} catch (Exception e) {
throw new IOException(e);
}
}
private void jsonRpc(HttpExchange x, String body) throws Exception {
ServerCallContext context = new ServerCallContext(UnauthenticatedUser.INSTANCE,
new HashMap<>(java.util.Map.of(ServerCallContext.TRANSPORT_KEY, TransportProtocol.JSONRPC)), Set.of());
A2AResponse<?> response = null;
Flow.Publisher<? extends A2AResponse<?>> stream = null;
try {
A2ARequest<?> request = JSONRPCUtils.parseRequestBody(body, "");
if (request instanceof SendMessageRequest r) {
response = rpc.onMessageSend(r, context);
} else if (request instanceof GetTaskRequest r) {
response = rpc.onGetTask(r, context);
} else if (request instanceof ListTasksRequest r) {
response = rpc.onListTasks(r, context);
} else if (request instanceof CancelTaskRequest r) {
response = rpc.onCancelTask(r, context);
} else if (request instanceof SendStreamingMessageRequest r) {
stream = rpc.onMessageSendStream(r, context);
} else if (request instanceof SubscribeToTaskRequest r) {
stream = rpc.onSubscribeToTask(r, context);
} else {
response = new A2AErrorResponse(request.getId(), new io.a2a.spec.UnsupportedOperationError());
}
} catch (A2AError e) {
response = new A2AErrorResponse(e);
} catch (Exception e) {
response = new A2AErrorResponse(new InternalError(e.getMessage()));
}
if (stream != null) {
sse(x, stream);
} else {
reply(x, serialize(response));
}
}
private static String serialize(A2AResponse<?> response) {
if (response instanceof A2AErrorResponse error) {
return JSONRPCUtils.toJsonRPCErrorResponse(error.getId(), error.getError());
}
if (response.getError() != null) {
return JSONRPCUtils.toJsonRPCErrorResponse(response.getId(), response.getError());
}
com.google.protobuf.MessageOrBuilder proto;
if (response instanceof GetTaskResponse r) {
proto = ProtoUtils.ToProto.task(r.getResult());
} else if (response instanceof CancelTaskResponse r) {
proto = ProtoUtils.ToProto.task(r.getResult());
} else if (response instanceof SendMessageResponse r) {
proto = ProtoUtils.ToProto.taskOrMessage(r.getResult());
} else if (response instanceof ListTasksResponse r) {
proto = ProtoUtils.ToProto.listTasksResult(r.getResult());
} else {
throw new IllegalArgumentException("Unsupported response: " + response.getClass().getName());
}
return JSONRPCUtils.toJsonRPCResultResponse(response.getId(), proto);
}
private void sse(HttpExchange x, Flow.Publisher<? extends A2AResponse<?>> publisher) throws IOException, InterruptedException {
x.getResponseHeaders().add("Content-Type", "text/event-stream");
x.sendResponseHeaders(200, 0);
OutputStream out = x.getResponseBody();
CountDownLatch done = new CountDownLatch(1);
AtomicLong id = new AtomicLong();
publisher.subscribe(new Flow.Subscriber<A2AResponse<?>>() {
@Override
public void onSubscribe(Flow.Subscription s) {
s.request(Long.MAX_VALUE);
}
@Override
public void onNext(A2AResponse<?> item) {
try {
out.write(SseFormatter.formatResponseAsSSE(item, id.getAndIncrement()).getBytes(StandardCharsets.UTF_8));
out.flush();
} catch (IOException e) {
done.countDown();
}
}
@Override
public void onError(Throwable t) {
done.countDown();
}
@Override
public void onComplete() {
done.countDown();
}
});
done.await();
}
private static void reply(HttpExchange x, String json) throws IOException {
byte[] bytes = json.getBytes(StandardCharsets.UTF_8);
x.getResponseHeaders().add("Content-Type", "application/json");
x.sendResponseHeaders(200, bytes.length);
x.getResponseBody().write(bytes);
}
/**
* The SDK starts its event-bus thread from a CDI lifecycle callback, and the public ensureStarted() does nothing
* outside CDI. Without this thread every event is accepted and none is delivered, and the client gets
* "Could not find a Task/Message". start() and stop() are package-private, so this calls them by reflection.
*/
private static void callPackagePrivate(MainEventBusProcessor processor, String method) {
try {
var m = MainEventBusProcessor.class.getDeclaredMethod(method);
m.setAccessible(true);
m.invoke(processor);
} catch (ReflectiveOperationException e) {
throw new IllegalStateException(e);
}
}
@Override
public void close() {
callPackagePrivate(processor, "stop");
http.stop(0);
executor.shutdownNow();
}
}
@@ -0,0 +1,79 @@
package com.ankurm.agents.a2a;
import io.a2a.server.agentexecution.AgentExecutor;
import io.a2a.server.agentexecution.RequestContext;
import io.a2a.server.tasks.AgentEmitter;
import io.a2a.spec.A2AError;
import io.a2a.spec.AgentCapabilities;
import io.a2a.spec.AgentCard;
import io.a2a.spec.AgentInterface;
import io.a2a.spec.AgentSkill;
import io.a2a.spec.Part;
import io.a2a.spec.TaskState;
import io.a2a.spec.TextPart;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CopyOnWriteArrayList;
/**
* Agent B: turns a comma-separated list of terms into a short brief. It does not know the definitions. For each
* term it asks Agent A (the glossary) over A2A, so one agent is the server for its caller and the client of another.
*/
public final class BriefAgent implements AgentExecutor {
private final String glossaryUrl;
/** What this agent sent to the glossary and what came back, for the transcript. */
public final List<String> callLog = new CopyOnWriteArrayList<>();
public BriefAgent(String glossaryUrl) {
this.glossaryUrl = glossaryUrl;
}
public static AgentCard card(String url) {
return AgentCard.builder()
.name("Brief Agent")
.description("Writes a short brief from a list of Java terms by asking the Glossary Agent")
.version("1.0.0")
.supportedInterfaces(List.of(new AgentInterface("JSONRPC", url)))
.capabilities(AgentCapabilities.builder().streaming(true).build())
.defaultInputModes(List.of("text/plain"))
.defaultOutputModes(List.of("text/plain"))
.skills(List.of(AgentSkill.builder()
.id("brief").name("Brief on several terms").description("Comma-separated terms in, bullet list out")
.tags(List.of("java", "summary")).examples(List.of("record, sealed class")).build()))
.build();
}
@Override
public void execute(RequestContext context, AgentEmitter emitter) throws A2AError {
emitter.submit();
emitter.startWork();
List<String> lines = new ArrayList<>();
for (String term : context.getUserInput(" ").split(",")) {
term = term.strip();
try {
A2aCaller.Reply reply = A2aCaller.send(glossaryUrl, term, null, false);
TaskState state = reply.task().status().state();
callLog.add("SendMessage \"" + term + "\" -> " + A2aCaller.short_(state) + " \"" + reply.text() + "\"");
if (state != TaskState.TASK_STATE_COMPLETED) {
emitter.fail(emitter.newAgentMessage(
List.<Part<?>>of(new TextPart("Glossary Agent could not help: " + reply.text())), Map.of()));
return;
}
lines.add("- " + term + ": " + reply.text());
} catch (Exception e) {
emitter.fail(emitter.newAgentMessage(List.<Part<?>>of(new TextPart("Glossary Agent unreachable: " + e.getMessage())), Map.of()));
return;
}
}
emitter.addArtifact(List.<Part<?>>of(new TextPart(String.join("\n", lines))), "brief", "brief", Map.of(), false, true);
emitter.complete();
}
@Override
public void cancel(RequestContext context, AgentEmitter emitter) throws A2AError {
emitter.cancel();
}
}
@@ -0,0 +1,70 @@
package com.ankurm.agents.a2a;
import io.a2a.server.agentexecution.AgentExecutor;
import io.a2a.server.agentexecution.RequestContext;
import io.a2a.server.tasks.AgentEmitter;
import io.a2a.spec.A2AError;
import io.a2a.spec.AgentCapabilities;
import io.a2a.spec.AgentCard;
import io.a2a.spec.AgentInterface;
import io.a2a.spec.AgentSkill;
import io.a2a.spec.Part;
import io.a2a.spec.TextPart;
import java.util.List;
import java.util.Map;
/**
* Agent A: defines Java terms. It does no model call, because the point is the protocol. A real agent
* would put an LLM or a database behind execute().
*/
public final class GlossaryAgent implements AgentExecutor {
static final Map<String, String> TERMS = Map.of(
"record", "A record is a final class that is a transparent carrier for its components.",
"sealed class", "A sealed class restricts which other classes may extend it.",
"virtual thread", "A virtual thread is a lightweight thread scheduled by the JVM, not the OS.");
public static AgentCard card(String url) {
return AgentCard.builder()
.name("Glossary Agent")
.description("Defines Java terms")
.version("1.0.0")
.supportedInterfaces(List.of(new AgentInterface("JSONRPC", url)))
.capabilities(AgentCapabilities.builder().streaming(true).build())
.defaultInputModes(List.of("text/plain"))
.defaultOutputModes(List.of("text/plain"))
.skills(List.of(AgentSkill.builder()
.id("define").name("Define a term").description("Returns a one-sentence definition")
.tags(List.of("java", "glossary")).examples(List.of("record")).build()))
.build();
}
@Override
public void execute(RequestContext context, AgentEmitter emitter) throws A2AError {
if (context.getTask() == null) {
emitter.submit();
}
String term = context.getUserInput(" ").strip().toLowerCase();
if (term.isEmpty()) {
emitter.requiresInput(emitter.newAgentMessage(List.<Part<?>>of(new TextPart("Which term should I define?")), Map.of()));
return;
}
emitter.startWork();
String definition = TERMS.get(term);
if (definition == null) {
emitter.fail(emitter.newAgentMessage(List.<Part<?>>of(new TextPart("I do not know the term '" + term + "'.")), Map.of()));
return;
}
// Two chunks of one artifact, so a streaming client sees the answer arrive in pieces.
int cut = definition.indexOf(' ', definition.length() / 2);
emitter.addArtifact(List.<Part<?>>of(new TextPart(definition.substring(0, cut))), "definition", "definition", Map.of(), false, false);
emitter.addArtifact(List.<Part<?>>of(new TextPart(definition.substring(cut))), "definition", "definition", Map.of(), true, true);
emitter.complete();
}
@Override
public void cancel(RequestContext context, AgentEmitter emitter) throws A2AError {
emitter.cancel();
}
}
@@ -0,0 +1,16 @@
package com.ankurm.agents.a2a;
import io.a2a.server.TransportMetadata;
import io.a2a.spec.TransportProtocol;
/**
* Tells the SDK that this application serves JSON-RPC. The SDK checks an agent card against the transports it
* can discover, and the Quarkus module normally registers one; since this demo has no Quarkus, it registers its own.
*/
public class JdkJsonRpcTransportMetadata implements TransportMetadata {
@Override
public String getTransportProtocol() {
return TransportProtocol.JSONRPC.asString();
}
}
@@ -0,0 +1 @@
com.ankurm.agents.a2a.JdkJsonRpcTransportMetadata
@@ -0,0 +1,203 @@
package com.ankurm.agents.a2a;
import com.google.gson.GsonBuilder;
import com.google.gson.JsonParser;
import io.a2a.spec.Task;
import org.junit.jupiter.api.Test;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import static org.assertj.core.api.Assertions.assertThat;
class A2aTest {
private static String get(String url) throws Exception {
return HttpClient.newHttpClient().send(HttpRequest.newBuilder(URI.create(url)).build(),
HttpResponse.BodyHandlers.ofString()).body();
}
private static String post(String url, String json) throws Exception {
return HttpClient.newHttpClient().send(HttpRequest.newBuilder(URI.create(url))
.header("Content-Type", "application/json").POST(HttpRequest.BodyPublishers.ofString(json)).build(),
HttpResponse.BodyHandlers.ofString()).body();
}
private static String text(Task t) {
StringBuilder sb = new StringBuilder();
t.artifacts().forEach(a -> a.parts().forEach(p -> sb.append(((io.a2a.spec.TextPart) p).text())));
return sb.toString();
}
private static String pretty(String json) {
return new GsonBuilder().setPrettyPrinting().create().toJson(JsonParser.parseString(json));
}
// ------------------------------------------------------------------ 1. the agent card
@Test
void theAgentCardIsPublishedAtAWellKnownPath() throws Exception {
try (A2aHttpServer glossary = A2aHttpServer.start(GlossaryAgent::card, new GlossaryAgent())) {
String card = pretty(get(glossary.url() + "/.well-known/agent-card.json"));
Transcript.write("01-agent-card.txt", "GET /.well-known/agent-card.json\n\n" + card + "\n");
assertThat(card).contains("\"name\": \"Glossary Agent\"").contains("\"streaming\": true").contains("\"id\": \"define\"");
}
}
// ------------------------------------------------------------------ 2. the raw exchange
@Test
void aMessageIsOneJsonRpcCallAndATaskIsTheAnswer() throws Exception {
try (A2aHttpServer glossary = A2aHttpServer.start(GlossaryAgent::card, new GlossaryAgent())) {
String request = "{\"jsonrpc\":\"2.0\",\"id\":\"1\",\"method\":\"SendMessage\",\"params\":{\"message\":"
+ "{\"messageId\":\"m-1\",\"role\":\"ROLE_USER\",\"parts\":[{\"text\":\"record\"}]},"
+ "\"configuration\":{\"blocking\":true}}}";
String response = post(glossary.url() + "/", request);
Transcript.write("02-raw-send-message.txt", "POST /\n" + pretty(request) + "\n\nResponse:\n" + pretty(response) + "\n");
assertThat(pretty(response)).contains("TASK_STATE_COMPLETED").contains("\"id\": 1,").contains("A record is a final class that is a transparent");
}
}
// ------------------------------------------------------------------ 3. blocking, streaming, and the task afterwards
@Test
void streamingShowsTheLifecycleAndTheTaskOutlivesTheCall() throws Exception {
try (A2aHttpServer glossary = A2aHttpServer.start(GlossaryAgent::card, new GlossaryAgent())) {
A2aCaller.Reply blocking = A2aCaller.send(glossary.url(), "record", null, false);
A2aCaller.Reply streamed = A2aCaller.send(glossary.url(), "record", null, true);
Task later = A2aCaller.getTask(glossary.url(), streamed.task().id());
StringBuilder sb = new StringBuilder("Same question, \"record\", asked twice.\n\n");
sb.append("Blocking (streaming off): ").append(blocking.events().size()).append(" event\n");
blocking.events().forEach(e -> sb.append(" ").append(A2aCaller.describe(e)).append('\n'));
sb.append(" answer: ").append(blocking.text()).append("\n\n");
sb.append("Streaming: ").append(streamed.events().size()).append(" events\n");
streamed.events().forEach(e -> sb.append(" ").append(A2aCaller.describe(e)).append('\n'));
sb.append(" answer: ").append(streamed.text()).append("\n\n");
sb.append("GetTask for the streamed task afterwards: state ").append(later.status().state())
.append(", artifacts ").append(later.artifacts().size()).append(", artifact text \"")
.append(text(later)).append("\"\n");
Transcript.write("03-streaming.txt", sb.toString());
assertThat(blocking.events()).hasSize(1);
assertThat(streamed.events()).hasSize(5);
assertThat(blocking.text()).isEqualTo(streamed.text());
assertThat(later.status().state().name()).isEqualTo("TASK_STATE_COMPLETED");
assertThat(text(later)).isEqualTo(streamed.text());
}
}
// ------------------------------------------------------------------ 4. input-required and a follow-up
@Test
void anAgentCanAskForMoreAndTheSameTaskContinues() throws Exception {
try (A2aHttpServer glossary = A2aHttpServer.start(GlossaryAgent::card, new GlossaryAgent())) {
A2aCaller.Reply first = A2aCaller.send(glossary.url(), " ", null, false);
String taskId = first.task().id();
A2aCaller.Reply second = A2aCaller.send(glossary.url(), "sealed class", taskId, false);
String sb = "Turn 1: message with no term.\n task " + taskId + " state " + A2aCaller.short_(first.task().status().state())
+ "\n agent says: " + first.text()
+ "\n\nTurn 2: message \"sealed class\" sent with taskId " + taskId + ".\n task " + second.task().id()
+ " state " + A2aCaller.short_(second.task().status().state()) + "\n answer: " + second.text()
+ "\n\nSame task id both turns: " + taskId.equals(second.task().id()) + "\n";
Transcript.write("04-input-required.txt", sb);
assertThat(first.task().status().state().name()).isEqualTo("TASK_STATE_INPUT_REQUIRED");
assertThat(second.task().id()).isEqualTo(taskId);
assertThat(second.text()).contains("restricts which other classes may extend it");
}
}
// ------------------------------------------------------------------ 5. two Java agents
@Test
void oneAgentCallsAnotherOverA2a() throws Exception {
try (A2aHttpServer glossary = A2aHttpServer.start(GlossaryAgent::card, new GlossaryAgent())) {
BriefAgent briefAgent = new BriefAgent(glossary.url());
try (A2aHttpServer brief = A2aHttpServer.start(BriefAgent::card, briefAgent)) {
A2aCaller.Reply ok = A2aCaller.send(brief.url(), "record, sealed class", null, false);
int callsAfterOk = briefAgent.callLog.size();
A2aCaller.Reply bad = A2aCaller.send(brief.url(), "record, monad", null, false);
StringBuilder sb = new StringBuilder("Caller -> Brief Agent -> Glossary Agent.\n\n");
sb.append("Request: \"record, sealed class\"\n Brief Agent task state: ").append(A2aCaller.short_(ok.task().status().state()))
.append("\n result:\n");
ok.text().lines().forEach(l -> sb.append(" ").append(l).append('\n'));
sb.append(" Calls the Brief Agent made to the Glossary Agent:\n");
briefAgent.callLog.subList(0, callsAfterOk).forEach(l -> sb.append(" ").append(l).append('\n'));
sb.append("\nRequest: \"record, monad\"\n Brief Agent task state: ").append(A2aCaller.short_(bad.task().status().state()))
.append("\n message: ").append(bad.text()).append("\n Calls the Brief Agent made to the Glossary Agent:\n");
briefAgent.callLog.subList(callsAfterOk, briefAgent.callLog.size()).forEach(l -> sb.append(" ").append(l).append('\n'));
Transcript.write("05-agent-to-agent.txt", sb.toString());
assertThat(ok.task().status().state().name()).isEqualTo("TASK_STATE_COMPLETED");
assertThat(ok.text()).contains("- record: A record is a final class").contains("- sealed class: A sealed class restricts");
assertThat(bad.task().status().state().name()).isEqualTo("TASK_STATE_FAILED");
assertThat(bad.text()).contains("monad");
}
}
}
// ------------------------------------------------------------------ 2b. without blocking the call returns early
@Test
void withoutBlockingTheCallReturnsBeforeTheAgentIsDone() throws Exception {
io.a2a.server.agentexecution.AgentExecutor slow = new io.a2a.server.agentexecution.AgentExecutor() {
private final GlossaryAgent inner = new GlossaryAgent();
@Override
public void execute(io.a2a.server.agentexecution.RequestContext c, io.a2a.server.tasks.AgentEmitter e) throws io.a2a.spec.A2AError {
e.submit();
e.startWork();
try {
Thread.sleep(400);
} catch (InterruptedException ex) {
Thread.currentThread().interrupt();
}
e.addArtifact(java.util.List.<io.a2a.spec.Part<?>>of(new io.a2a.spec.TextPart("slow answer")), "answer", "answer", java.util.Map.of(), false, true);
e.complete();
}
@Override
public void cancel(io.a2a.server.agentexecution.RequestContext c, io.a2a.server.tasks.AgentEmitter e) throws io.a2a.spec.A2AError {
inner.cancel(c, e);
}
};
try (A2aHttpServer glossary = A2aHttpServer.start(GlossaryAgent::card, slow)) {
String body = "{\"jsonrpc\":\"2.0\",\"id\":\"1\",\"method\":\"SendMessage\",\"params\":{\"message\":"
+ "{\"messageId\":\"m-%s\",\"role\":\"ROLE_USER\",\"parts\":[{\"text\":\"record\"}]}%s}}";
long t0 = System.nanoTime();
String early = post(glossary.url() + "/", body.formatted("a", ",\"configuration\":{\"blocking\":false}"));
long earlyMs = (System.nanoTime() - t0) / 1_000_000;
t0 = System.nanoTime();
String blocking = post(glossary.url() + "/", body.formatted("b", ",\"configuration\":{\"blocking\":true}"));
long blockingMs = (System.nanoTime() - t0) / 1_000_000;
t0 = System.nanoTime();
String omitted = post(glossary.url() + "/", body.formatted("c", ""));
long omittedMs = (System.nanoTime() - t0) / 1_000_000;
String omittedState = com.google.gson.JsonParser.parseString(omitted).getAsJsonObject().getAsJsonObject("result")
.getAsJsonObject("task").getAsJsonObject("status").get("state").getAsString();
String earlyState = com.google.gson.JsonParser.parseString(early).getAsJsonObject().getAsJsonObject("result")
.getAsJsonObject("task").getAsJsonObject("status").get("state").getAsString();
String blockingState = com.google.gson.JsonParser.parseString(blocking).getAsJsonObject().getAsJsonObject("result")
.getAsJsonObject("task").getAsJsonObject("status").get("state").getAsString();
String sb = "The agent accepts the task at once, then works for 400 ms before finishing.\n\n"
+ "SendMessage with configuration.blocking = false: task state " + earlyState + ", returned in under 400 ms: "
+ (earlyMs < 400) + "\n"
+ "SendMessage with configuration.blocking = true: task state " + blockingState + ", returned in at least 400 ms: "
+ (blockingMs >= 400) + "\n"
+ "SendMessage with no configuration at all: task state " + omittedState + ", returned in under 400 ms: "
+ (omittedMs < 400) + "\n";
Transcript.write("07-blocking.txt", sb);
assertThat(earlyState).isIn("TASK_STATE_SUBMITTED", "TASK_STATE_WORKING");
assertThat(blockingState).isEqualTo("TASK_STATE_COMPLETED");
assertThat(omittedState).isIn("TASK_STATE_SUBMITTED", "TASK_STATE_WORKING");
assertThat(omittedMs).isLessThan(400);
assertThat(earlyMs).isLessThan(400);
assertThat(blockingMs).isGreaterThanOrEqualTo(400);
}
}
}
@@ -0,0 +1,39 @@
package com.ankurm.agents.a2a;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.regex.Pattern;
/** Writes what a test observed to output/NN-name.txt so every figure in the post comes from a file. */
final class Transcript {
private static final Path DIR = Path.of("output");
private static final Pattern UUID = Pattern.compile("[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}");
private static final Pattern PORT = Pattern.compile("127\\.0\\.0\\.1:\\d+");
private static final Pattern TS = Pattern.compile("\\d{4}-\\d\\d-\\d\\dT[\\d:.]+Z");
private Transcript() {
}
/** Ports, ids and timestamps change on every run, so they are replaced by stable placeholders. */
static String scrub(String s) {
var ids = new java.util.LinkedHashMap<String, String>();
var m = UUID.matcher(s);
StringBuilder out = new StringBuilder();
while (m.find()) {
m.appendReplacement(out, ids.computeIfAbsent(m.group(), k -> "<id-" + (ids.size() + 1) + ">"));
}
m.appendTail(out);
return TS.matcher(PORT.matcher(out.toString()).replaceAll("127.0.0.1:PORT")).replaceAll("<time>");
}
static void write(String name, String content) {
try {
Files.createDirectories(DIR);
Files.writeString(DIR.resolve(name), scrub(content));
} catch (IOException e) {
throw new IllegalStateException(e);
}
}
}
@@ -0,0 +1 @@
org.slf4j.simpleLogger.defaultLogLevel=warn
+2
View File
@@ -11,6 +11,7 @@
<modules>
<module>embabel</module>
<module>a2a</module>
</modules>
<properties>
@@ -18,6 +19,7 @@
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<!-- Verified on Maven Central 2026-10-11 -->
<embabel.version>1.5.3</embabel.version>
<a2a.sdk.version>1.0.0.Alpha3</a2a.sdk.version>
<junit.version>6.1.3</junit.version>
<assertj.version>3.27.6</assertj.version>
<slf4j.version>2.0.17</slf4j.version>