Skip to content

Commit 56d0431

Browse files
author
Mark Pollack
committed
Reconnect a dropped SSE stream instead of failing the transport
The Streamable HTTP client ended the whole transport when its connection stream closed, and when a session stream with a pending response closed, even when the server had only detached it (backpressure, a proxy in between restarting). The server keeps undelivered events in the stream's mailbox, so the client now reopens the stream with a short backoff and carries on. It still gives up when the server answers 404, meaning the connection or session is gone, or after three reconnects in a row that delivered no event, so a server that keeps closing cannot cause an endless loop. The reopened stream replaces the old one only if nothing else already did.
1 parent 1b63a9a commit 56d0431

3 files changed

Lines changed: 131 additions & 8 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,11 @@ building a second client on an already-connected transport now fails at construc
3939
GET first on `http://` endpoints, so against a server that speaks h2c (the SDK's own does) every
4040
request, streams included, runs on HTTP/2. Against one that does not, the probe settles on HTTP/1.1
4141
and later requests stop offering the upgrade, which some servers route to their WebSocket handler.
42+
- **The Streamable HTTP client reconnects a dropped SSE stream.** When the server or the network closes
43+
the connection stream or a session stream, the client reopens it with a short backoff instead of
44+
failing the transport; the server's mailbox delivers whatever was not yet written. It gives up, as
45+
before, when the server answers 404 (the connection is gone) or after three reconnects in a row that
46+
delivered nothing.
4247
- **Client sessions learn that their transport died.** `AcpClientTransport.awaitTermination()` (default:
4348
never) is implemented by the Streamable HTTP and WebSocket client transports; `AcpClientSession` fails
4449
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: 63 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -783,22 +783,73 @@ private void clearState() {
783783
}
784784
}
785785

786+
/** Reconnects in a row that delivered nothing before the transport gives up on a stream. */
787+
private static final int MAX_BARREN_RECONNECTS = 3;
788+
789+
private static final Duration RECONNECT_BACKOFF = Duration.ofMillis(200);
790+
791+
/**
792+
* A stream closed that this client did not close: the server detached it (backpressure,
793+
* a restart of the proxy in between), or the network dropped it. The server keeps what
794+
* it has not delivered in the stream's mailbox, so reopening loses nothing: reconnect
795+
* with a short backoff. Give up, as before, when the server answers that the connection
796+
* is gone, or after {@value #MAX_BARREN_RECONNECTS} reconnects in a row that delivered no
797+
* event (a server that keeps closing must not cause an endless loop).
798+
*/
786799
private void handleUnexpectedSseClosure(SseStream stream, Throwable error) {
787800
if (closing.get()) {
788801
return;
789802
}
790-
791803
RouteScope scope = stream.scope;
792-
if (!scope.isSession()) {
793-
terminateAfterSseFailure(error);
804+
int barren = stream.delivered ? 0 : stream.barrenReconnects + 1;
805+
if (barren > MAX_BARREN_RECONNECTS) {
806+
giveUpOn(stream, error);
794807
return;
795808
}
809+
logger.info("SSE stream closed unexpectedly; reconnecting: {}", scope);
810+
Mono.defer(() -> openSseStream(scope))
811+
.retryWhen(reactor.util.retry.Retry.backoff(2, RECONNECT_BACKOFF)
812+
.scheduler(AcpSchedulers.timeouts())
813+
.filter(e -> !isConnectionGone(e)))
814+
.subscribe(reopened -> {
815+
reopened.barrenReconnects = barren;
816+
if (closing.get() || !replaceStream(stream, reopened)) {
817+
reopened.close();
818+
return;
819+
}
820+
reopened.start();
821+
logger.info("SSE stream reconnected: {}", scope);
822+
}, reconnectError -> giveUpOn(stream, reconnectError));
823+
}
824+
825+
private boolean replaceStream(SseStream old, SseStream reopened) {
826+
if (old.scope.isSession()) {
827+
return sessionStreams.replace(old.scope.sessionId(), old, reopened);
828+
}
829+
synchronized (this) {
830+
if (this.connectionStream != old) {
831+
return false;
832+
}
833+
this.connectionStream = reopened;
834+
return true;
835+
}
836+
}
837+
838+
/** 404 on reconnect: the server no longer knows this connection or session. */
839+
private static boolean isConnectionGone(Throwable error) {
840+
return error instanceof AcpConnectionException && String.valueOf(error.getMessage()).contains("got 404");
841+
}
796842

797-
if (hasPendingResponseFor(scope)) {
843+
/** The old behaviour: a dead connection stream, or a session stream owing a response, ends the transport. */
844+
private void giveUpOn(SseStream stream, Throwable error) {
845+
if (closing.get()) {
846+
return;
847+
}
848+
RouteScope scope = stream.scope;
849+
if (!scope.isSession() || hasPendingResponseFor(scope)) {
798850
terminateAfterSseFailure(error);
799851
return;
800852
}
801-
802853
if (sessionStreams.remove(scope.sessionId(), stream)) {
803854
sessionStreamOpenOperations.remove(scope.sessionId());
804855
logger.info("Session SSE stream closed; it will be reopened before the next session request: {}", scope);
@@ -873,6 +924,12 @@ private class SseStream {
873924

874925
private Future<?> readerTask;
875926

927+
/** Whether this stream carried at least one event; resets the reconnect budget. */
928+
private volatile boolean delivered;
929+
930+
/** Reconnects in a row, ending with this stream, that delivered nothing. */
931+
private volatile int barrenReconnects;
932+
876933
SseStream(RouteScope scope, InputStream body) {
877934
this.scope = scope;
878935
this.body = body;
@@ -939,6 +996,7 @@ private void dispatchEvent(StringBuilder dataBuffer) {
939996
}
940997
try {
941998
JSONRPCMessage message = AcpSchema.deserializeJsonRpcMessage(jsonMapper, dataBuffer.toString());
999+
delivered = true;
9421000
processInbound(scope, message).subscribe(v -> {
9431001
}, error -> {
9441002
if (!closed.get() && !closing.get()) {

‎acp-core/src/test/java/com/agentclientprotocol/sdk/client/transport/StreamableHttpAcpClientTransportTest.java‎

Lines changed: 63 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -828,9 +828,12 @@ void cleartextServerWithoutH2cGetsPlainHttp11Requests() throws Exception {
828828
transport.sendMessage(AcpTestFixtures.createJsonRpcRequest(AcpSchema.METHOD_INITIALIZE, "init-1",
829829
AcpTestFixtures.createInitializeRequest())).block();
830830

831-
assertThat(requests.get(0).method()).as("the probe").isEqualTo("GET");
832-
assertThat(requests.get(0).headers().firstValue("Acp-Connection-Id")).isEmpty();
833-
assertThat(requests.subList(1, requests.size()))
831+
// The mocked connection stream ends at once, so the client keeps reconnecting while
832+
// this runs; take a snapshot. Reconnects must be pinned too.
833+
List<HttpRequest> seen = List.copyOf(requests);
834+
assertThat(seen.get(0).method()).as("the probe").isEqualTo("GET");
835+
assertThat(seen.get(0).headers().firstValue("Acp-Connection-Id")).isEmpty();
836+
assertThat(seen.subList(1, seen.size()))
834837
.as("every request after the probe is pinned to HTTP/1.1")
835838
.isNotEmpty()
836839
.allMatch(request -> request.version().equals(java.util.Optional.of(HttpClient.Version.HTTP_1_1)));
@@ -858,6 +861,63 @@ void initializeRequiresConnectionIdHeader() throws Exception {
858861
.hasMessageContaining("Acp-Connection-Id");
859862
}
860863

864+
/**
865+
* The server closed the connection stream once (backpressure, a proxy restart). The
866+
* client reconnects, and a response the server sends afterwards is delivered: its
867+
* mailbox kept it.
868+
*/
869+
@Test
870+
void droppedConnectionStreamIsReconnectedAndLaterResponsesArrive() throws Exception {
871+
HttpClient httpClient = mock(HttpClient.class);
872+
AtomicInteger connectionGets = new AtomicInteger();
873+
PipedInputStream secondBody = new PipedInputStream();
874+
PipedOutputStream secondWriter = new PipedOutputStream(secondBody);
875+
BlockingQueue<AcpSchema.JSONRPCMessage> inbound = new LinkedBlockingQueue<>();
876+
when(httpClient.sendAsync(any(), any())).thenAnswer(invocation -> {
877+
HttpRequest request = invocation.getArgument(0);
878+
if ("POST".equals(request.method()) && request.headers().firstValue("Acp-Connection-Id").isEmpty()) {
879+
String init = jsonMapper.writeValueAsString(AcpTestFixtures.createJsonRpcResponse("init-1",
880+
AcpTestFixtures.createInitializeResponse()));
881+
return CompletableFuture.completedFuture(response(200,
882+
Map.of("Content-Type", "application/json", "Acp-Connection-Id", "conn-1"), init));
883+
}
884+
if ("GET".equals(request.method()) && request.headers().firstValue("Acp-Session-Id").isEmpty()) {
885+
// First connection stream ends at once; the reconnect gets a live one.
886+
InputStream body = connectionGets.incrementAndGet() == 1 ? emptyBody() : secondBody;
887+
return CompletableFuture.completedFuture(response(200, Map.of("Content-Type", "text/event-stream"), body));
888+
}
889+
if ("GET".equals(request.method())) {
890+
return CompletableFuture.completedFuture(
891+
response(200, Map.of("Content-Type", "text/event-stream"), new PipedInputStream(1024)));
892+
}
893+
return CompletableFuture.completedFuture(response(202, Map.of(), null));
894+
});
895+
StreamableHttpAcpClientTransport transport = new StreamableHttpAcpClientTransport(
896+
URI.create("https://localhost:8443/acp"), jsonMapper, httpClient);
897+
try {
898+
transport.connect(message -> message.doOnNext(inbound::add).then(Mono.empty())).block();
899+
transport.sendMessage(AcpTestFixtures.createJsonRpcRequest(AcpSchema.METHOD_INITIALIZE, "init-1",
900+
AcpTestFixtures.createInitializeRequest())).block();
901+
awaitResponse(inbound, "init-1");
902+
903+
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
904+
while (connectionGets.get() < 2 && System.nanoTime() < deadline) {
905+
Thread.sleep(20);
906+
}
907+
assertThat(connectionGets.get()).as("the connection stream was reopened").isGreaterThanOrEqualTo(2);
908+
909+
transport.sendMessage(AcpTestFixtures.createJsonRpcRequest(AcpSchema.METHOD_SESSION_NEW, "new-1",
910+
AcpTestFixtures.createNewSessionRequest("/workspace"))).block();
911+
writeSse(secondWriter, new AcpSchema.JSONRPCResponse(AcpSchema.JSONRPC_VERSION, "new-1",
912+
new AcpSchema.NewSessionResponse("sess-1", null, null), null));
913+
assertThat(awaitResponse(inbound, "new-1").error()).isNull();
914+
}
915+
finally {
916+
secondWriter.close();
917+
transport.close();
918+
}
919+
}
920+
861921
private void awaitSessionSseClosure() throws InterruptedException {
862922
Thread.sleep(100);
863923
}

0 commit comments

Comments
 (0)