Skip to content

Commit 99c7e0f

Browse files
author
Mark Pollack
committed
Negotiate HTTP/2 from the first request over plain http
The Streamable HTTP transport requires HTTP/2, and localhost without TLS is a first-class deployment. Over cleartext the JDK HttpClient only offers the h2c upgrade on a request without a body, so initialize, a POST, went out on HTTP/1.1 and the connection upgraded only on the first SSE GET. The client now sends a bodiless OPTIONS before initialize on http:// endpoints; a failed probe is ignored. A new test drives the SDK client against the SDK server over http:// through a recording HttpClient and asserts every response, initialize, POSTs and the SSE streams, was HTTP/2; it fails without the probe.
1 parent 8322751 commit 99c7e0f

3 files changed

Lines changed: 156 additions & 1 deletion

File tree

‎CHANGELOG.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
6363
the SSE mailbox and per-subscriber queue limits are configurable, attached streams get a `: keep-alive`
6464
comment every 15 s so proxies do not cut idle connections, and a new GET on a stream takes it over
6565
from a subscriber the server may not yet know is dead instead of fanning out duplicates.
66+
- **HTTP/2 over plain `http://`.** The RFD requires HTTP/2, and localhost without TLS is a first-class
67+
deployment. Over cleartext the JDK client only offers the h2c upgrade on a request without a body, so
68+
`initialize`, a POST, went out on HTTP/1.1. `StreamableHttpAcpClientTransport` now sends a bodiless
69+
OPTIONS first on `http://` endpoints, and every request, streams included, runs on HTTP/2.
6670
- **Client sessions learn that their transport died.** `AcpClientTransport.awaitTermination()` (default:
6771
never) is implemented by the Streamable HTTP and WebSocket client transports; `AcpClientSession` fails
6872
pending requests at once with the cause, and every later request, instead of waiting out the request

‎acp-core/src/main/java/com/agentclientprotocol/sdk/client/transport/StreamableHttpAcpClientTransport.java‎

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -298,6 +298,30 @@ public Mono<Void> sendMessage(JSONRPCMessage message) {
298298
return routeAndPost(message);
299299
}
300300

301+
/**
302+
* The transport requires HTTP/2 (RFD). Over {@code https} ALPN negotiates it on the first
303+
* request. Over cleartext {@code http} the JDK client only offers the h2c upgrade on a
304+
* request without a body, and {@code initialize} is a POST, so without this the whole
305+
* connection would stay on HTTP/1.1. A bodiless OPTIONS first upgrades the connection;
306+
* every later request, the POSTs and the SSE GETs, reuses it over HTTP/2. A failed probe
307+
* is ignored: initialize then proceeds and reports any real problem itself.
308+
*/
309+
private Mono<Void> upgradeCleartextToHttp2() {
310+
if (!"http".equalsIgnoreCase(endpointUri.getScheme()) || httpClient.version() != HttpClient.Version.HTTP_2) {
311+
return Mono.empty();
312+
}
313+
HttpRequest probe = HttpRequest.newBuilder(endpointUri)
314+
.method("OPTIONS", HttpRequest.BodyPublishers.noBody())
315+
.build();
316+
return sendAsync(probe, HttpResponse.BodyHandlers.discarding())
317+
.doOnNext(response -> logger.debug("Cleartext probe to {} negotiated {}", endpointUri, response.version()))
318+
.then()
319+
.onErrorResume(error -> {
320+
logger.debug("Cleartext HTTP/2 probe to {} failed: {}", endpointUri, error.getMessage());
321+
return Mono.empty();
322+
});
323+
}
324+
301325
private Mono<Void> initialize(AcpSchema.JSONRPCRequest request) {
302326
if (!initialized.compareAndSet(false, true)) {
303327
return Mono.error(new IllegalStateException("Transport is already initialized"));
@@ -314,7 +338,8 @@ private Mono<Void> initialize(AcpSchema.JSONRPCRequest request) {
314338
return Mono.error(new AcpConnectionException("Failed to serialize initialize request", e));
315339
}
316340

317-
return sendAsync(httpRequest, HttpResponse.BodyHandlers.ofString())
341+
return upgradeCleartextToHttp2()
342+
.then(sendAsync(httpRequest, HttpResponse.BodyHandlers.ofString()))
318343
.flatMap(response -> {
319344
if (response.statusCode() != 200) {
320345
return Mono.error(new AcpConnectionException(

‎acp-streamable-http-jetty/src/test/java/com/agentclientprotocol/sdk/agent/transport/StreamableHttpClientServerIntegrationTest.java‎

Lines changed: 126 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -210,6 +210,132 @@ void initializeIsServedOverHttp2() throws Exception {
210210
}
211211
}
212212

213+
/**
214+
* Localhost over plain {@code http://} is a first-class deployment (no TLS), and the RFD
215+
* requires HTTP/2: every request the SDK client makes, initialize, the POSTs and the
216+
* long-lived SSE GETs, must run on HTTP/2 over cleartext.
217+
*/
218+
@Test
219+
void sdkClientUsesHttp2ForEveryRequestOverCleartextHttp() throws Exception {
220+
StreamableHttpAcpAgentTransport server = startServer(StreamableHttpAcpAgentTransportOptions.defaults());
221+
RecordingHttpClient http = new RecordingHttpClient(
222+
HttpClient.newBuilder().version(HttpClient.Version.HTTP_2).build());
223+
AcpAsyncClient client = AcpClient
224+
.async(new StreamableHttpAcpClientTransport(endpoint(server), AcpJsonMapper.createDefault(), http))
225+
.requestTimeout(TIMEOUT)
226+
.build();
227+
try {
228+
client.initialize().block(TIMEOUT);
229+
AcpSchema.NewSessionResponse session = client
230+
.newSession(new AcpSchema.NewSessionRequest("/workspace", List.of()))
231+
.block(TIMEOUT);
232+
client.prompt(new AcpSchema.PromptRequest(session.sessionId(), List.of(new AcpSchema.TextContent("hi"))))
233+
.block(TIMEOUT);
234+
235+
assertThat(endpoint(server).getScheme()).isEqualTo("http");
236+
assertThat(http.versions()).as("request methods seen: " + http.methods()).isNotEmpty();
237+
assertThat(http.methods()).contains("POST", "GET");
238+
assertThat(http.versions()).as("every response, including SSE streams").containsOnly(HttpClient.Version.HTTP_2);
239+
}
240+
finally {
241+
client.closeGracefully().block(TIMEOUT);
242+
server.closeGracefully().block(TIMEOUT);
243+
}
244+
}
245+
246+
/** Delegates to a real client and records the method and negotiated version of every response. */
247+
private static final class RecordingHttpClient extends HttpClient {
248+
249+
private final HttpClient delegate;
250+
251+
private final List<HttpClient.Version> versions = new CopyOnWriteArrayList<>();
252+
253+
private final List<String> methods = new CopyOnWriteArrayList<>();
254+
255+
RecordingHttpClient(HttpClient delegate) {
256+
this.delegate = delegate;
257+
}
258+
259+
List<HttpClient.Version> versions() {
260+
return versions;
261+
}
262+
263+
List<String> methods() {
264+
return methods;
265+
}
266+
267+
private <T> HttpResponse<T> record(HttpRequest request, HttpResponse<T> response) {
268+
methods.add(request.method());
269+
versions.add(response.version());
270+
return response;
271+
}
272+
273+
@Override
274+
public java.util.Optional<java.net.CookieHandler> cookieHandler() {
275+
return delegate.cookieHandler();
276+
}
277+
278+
@Override
279+
public java.util.Optional<Duration> connectTimeout() {
280+
return delegate.connectTimeout();
281+
}
282+
283+
@Override
284+
public Redirect followRedirects() {
285+
return delegate.followRedirects();
286+
}
287+
288+
@Override
289+
public java.util.Optional<java.net.ProxySelector> proxy() {
290+
return delegate.proxy();
291+
}
292+
293+
@Override
294+
public javax.net.ssl.SSLContext sslContext() {
295+
return delegate.sslContext();
296+
}
297+
298+
@Override
299+
public javax.net.ssl.SSLParameters sslParameters() {
300+
return delegate.sslParameters();
301+
}
302+
303+
@Override
304+
public java.util.Optional<java.net.Authenticator> authenticator() {
305+
return delegate.authenticator();
306+
}
307+
308+
@Override
309+
public Version version() {
310+
return delegate.version();
311+
}
312+
313+
@Override
314+
public java.util.Optional<java.util.concurrent.Executor> executor() {
315+
return delegate.executor();
316+
}
317+
318+
@Override
319+
public <T> HttpResponse<T> send(HttpRequest request, HttpResponse.BodyHandler<T> handler)
320+
throws IOException, InterruptedException {
321+
return record(request, delegate.send(request, handler));
322+
}
323+
324+
@Override
325+
public <T> java.util.concurrent.CompletableFuture<HttpResponse<T>> sendAsync(HttpRequest request,
326+
HttpResponse.BodyHandler<T> handler) {
327+
return delegate.sendAsync(request, handler).thenApply(response -> record(request, response));
328+
}
329+
330+
@Override
331+
public <T> java.util.concurrent.CompletableFuture<HttpResponse<T>> sendAsync(HttpRequest request,
332+
HttpResponse.BodyHandler<T> handler, HttpResponse.PushPromiseHandler<T> pushPromiseHandler) {
333+
return delegate.sendAsync(request, handler, pushPromiseHandler)
334+
.thenApply(response -> record(request, response));
335+
}
336+
337+
}
338+
213339
/** An agent factory that fails must fail the WebSocket upgrade, not leave a half-open connection. */
214340
@Test
215341
void webSocketUpgradeFailsCleanlyWhenTheAgentFactoryThrows() throws Exception {

0 commit comments

Comments
 (0)