From 81dfe0e2cd325bab530fb102638ac760b0183e28 Mon Sep 17 00:00:00 2001 From: AlgoVoi Date: Wed, 23 Sep 2026 15:52:15 +0000 Subject: [PATCH 1/3] fix(client): deliver exactly one terminal SSE callback (#1170) The JSON-RPC SSE listeners guarded only onComplete() with a volatile completed flag, so onError, a parse failure, and post-cancellation signals could each still reach the error/completion consumer, delivering more than one terminal callback (including a null completion after an error or cancellation). Introduce one shared atomic transition. AbstractSSEEventListener now owns an AtomicBoolean and a signalTerminal(Throwable) helper that lets the first caller win via compareAndSet and delivers exactly one outcome (a non-null failure or a null normal completion); every later signal is dropped. onError, onComplete and the parse-error path all route through it. The 0.3 compatibility JSON-RPC listener, which does not share the base class, mirrors the same AtomicBoolean pattern. REST is left as is (both native and 0.3 use a no-op completion callback, an API decision). Tests cover complete-then-error, error-then-complete, repeated completion, a 32-thread concurrent race, and a final event followed by completion, asserting exactly one terminal callback in each. Signed-off-by: AlgoVoi --- .../jsonrpc/sse/SSEEventListener.java | 21 +--- .../jsonrpc/sse/SSEEventListenerTest.java | 106 ++++++++++++++++++ .../spi/sse/AbstractSSEEventListener.java | 27 ++++- .../jsonrpc/sse/SSEEventListener_v0_3.java | 38 +++---- .../sse/SSEEventListener_v0_3_Test.java | 89 +++++++++++++++ 5 files changed, 237 insertions(+), 44 deletions(-) diff --git a/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListener.java b/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListener.java index 63d0015b8..978fcaf70 100644 --- a/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListener.java +++ b/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListener.java @@ -21,8 +21,6 @@ public class SSEEventListener extends AbstractSSEEventListener { private static final Logger LOGGER = Logger.getLogger(SSEEventListener.class.getName()); - private volatile boolean completed = false; - public SSEEventListener(Consumer eventHandler, @Nullable Consumer errorHandler) { super(eventHandler, errorHandler); @@ -34,21 +32,8 @@ public void onMessage(ServerSentEvent event, @Nullable Future completableF } public void onComplete() { - // Idempotent: only signal completion once, even if called multiple times - if (completed) { - LOGGER.fine("SSEEventListener.onComplete() called again - ignoring (already completed)"); - return; - } - completed = true; - - // Signal normal stream completion (null error means successful completion) LOGGER.fine("SSEEventListener.onComplete() called - signaling successful stream completion"); - if (getErrorHandler() != null) { - LOGGER.fine("Calling errorHandler.accept(null) to signal successful completion"); - getErrorHandler().accept(null); - } else { - LOGGER.warning("errorHandler is null, cannot signal completion"); - } + signalTerminal(null); } /** @@ -65,9 +50,7 @@ private void parseAndHandleMessage(String message, @Nullable Future future // Delegate to base class for common event handling and auto-close logic handleEvent(event, future); } catch (A2AError error) { - if (getErrorHandler() != null) { - getErrorHandler().accept(error); - } + signalTerminal(error); } catch (JsonProcessingException e) { throw new RuntimeException(e); } diff --git a/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListenerTest.java b/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListenerTest.java index a256f253c..a3424d99c 100644 --- a/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListenerTest.java +++ b/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListenerTest.java @@ -4,12 +4,18 @@ import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertTrue; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; @@ -248,6 +254,106 @@ public void testOnEventWithFinalTaskStatusUpdateEventEventCancels() throws Excep } + + @Test + public void testOnCompleteThenOnErrorDeliversSingleTerminalCallback() { + AtomicInteger terminalCount = new AtomicInteger(0); + AtomicReference lastArg = new AtomicReference<>(); + SSEEventListener listener = new SSEEventListener( + event -> {}, + error -> { terminalCount.incrementAndGet(); lastArg.set(error); }); + + listener.onComplete(); + listener.onError(new RuntimeException("late error"), new CancelCapturingFuture()); + + assertEquals(1, terminalCount.get()); + assertNull(lastArg.get()); + } + + @Test + public void testOnErrorThenOnCompleteDeliversSingleTerminalCallback() { + AtomicInteger terminalCount = new AtomicInteger(0); + AtomicReference lastArg = new AtomicReference<>(); + SSEEventListener listener = new SSEEventListener( + event -> {}, + error -> { terminalCount.incrementAndGet(); lastArg.set(error); }); + + RuntimeException boom = new RuntimeException("first error"); + listener.onError(boom, new CancelCapturingFuture()); + listener.onComplete(); + + assertEquals(1, terminalCount.get()); + assertSame(boom, lastArg.get()); + } + + @Test + public void testRepeatedOnCompleteDeliversSingleTerminalCallback() { + AtomicInteger terminalCount = new AtomicInteger(0); + AtomicReference lastArg = new AtomicReference<>(); + SSEEventListener listener = new SSEEventListener( + event -> {}, + error -> { terminalCount.incrementAndGet(); lastArg.set(error); }); + + listener.onComplete(); + listener.onComplete(); + listener.onComplete(); + + assertEquals(1, terminalCount.get()); + } + + @Test + public void testConcurrentTerminalSignalsDeliverExactlyOneCallback() throws Exception { + AtomicInteger terminalCount = new AtomicInteger(0); + AtomicReference lastArg = new AtomicReference<>(); + SSEEventListener listener = new SSEEventListener( + event -> {}, + error -> { terminalCount.incrementAndGet(); lastArg.set(error); }); + + int threads = 32; + ExecutorService pool = Executors.newFixedThreadPool(threads); + CountDownLatch start = new CountDownLatch(1); + CountDownLatch done = new CountDownLatch(threads); + for (int i = 0; i < threads; i++) { + final boolean complete = (i % 2 == 0); + pool.submit(() -> { + try { + start.await(); + if (complete) { + listener.onComplete(); + } else { + listener.onError(new RuntimeException("race"), new CancelCapturingFuture()); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + done.countDown(); + } + }); + } + start.countDown(); + assertTrue(done.await(10, TimeUnit.SECONDS)); + pool.shutdownNow(); + + assertEquals(1, terminalCount.get()); + } + + @Test + public void testFinalEventThenOnCompleteDeliversSingleTerminalCallback() throws Exception { + AtomicInteger terminalCount = new AtomicInteger(0); + SSEEventListener listener = new SSEEventListener( + event -> {}, + error -> terminalCount.incrementAndGet()); + + String eventData = JsonStreamingMessages.STREAMING_STATUS_UPDATE_EVENT_FINAL.substring( + JsonStreamingMessages.STREAMING_STATUS_UPDATE_EVENT_FINAL.indexOf("{")); + CancelCapturingFuture future = new CancelCapturingFuture(); + listener.onMessage(new ServerSentEvent(eventData), future); + listener.onComplete(); + + assertTrue(future.cancelHandlerCalled); + assertEquals(1, terminalCount.get()); + } + private static class CancelCapturingFuture implements Future { private boolean cancelHandlerCalled; diff --git a/client/transport/spi/src/main/java/org/a2aproject/sdk/client/transport/spi/sse/AbstractSSEEventListener.java b/client/transport/spi/src/main/java/org/a2aproject/sdk/client/transport/spi/sse/AbstractSSEEventListener.java index 8bb09a601..29b30d000 100644 --- a/client/transport/spi/src/main/java/org/a2aproject/sdk/client/transport/spi/sse/AbstractSSEEventListener.java +++ b/client/transport/spi/src/main/java/org/a2aproject/sdk/client/transport/spi/sse/AbstractSSEEventListener.java @@ -1,6 +1,7 @@ package org.a2aproject.sdk.client.transport.spi.sse; import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Consumer; import java.util.logging.Logger; @@ -25,6 +26,7 @@ public abstract class AbstractSSEEventListener { private final Consumer eventHandler; private final @Nullable Consumer errorHandler; + private final AtomicBoolean terminalSignaled = new AtomicBoolean(false); /** * Creates a new SSE event listener with the specified handlers. @@ -73,14 +75,33 @@ protected Consumer getEventHandler() { * @param future Optional future for closing the SSE connection */ public void onError(Throwable throwable, @Nullable Future future) { - if (errorHandler != null) { - errorHandler.accept(throwable); - } + signalTerminal(throwable); if (future != null) { future.cancel(true); // close SSE channel } } + /** + * Delivers exactly one terminal callback for the stream. The first caller to win + * the atomic transition delivers its outcome to the error/completion consumer (a + * non-null {@code error} is a failure, {@code null} is normal completion); every + * later completion, error or post-cancellation signal is dropped, so a single + * streaming request yields exactly one terminal callback. + * + * @param error the failure to report, or {@code null} to signal normal completion + */ + protected void signalTerminal(@Nullable Throwable error) { + if (!terminalSignaled.compareAndSet(false, true)) { + LOGGER.fine("Terminal callback already delivered, ignoring subsequent signal"); + return; + } + if (errorHandler != null) { + errorHandler.accept(error); + } else if (error != null) { + LOGGER.warning("errorHandler is null, cannot report terminal error"); + } + } + /** * Processes a parsed streaming event and handles auto-close logic for final events. * This method encapsulates the common logic for handling events and determining diff --git a/compat-0.3/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3.java b/compat-0.3/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3.java index 75ec20c46..22b18a4be 100644 --- a/compat-0.3/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3.java +++ b/compat-0.3/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3.java @@ -10,6 +10,7 @@ import org.a2aproject.sdk.compat03.spec.TaskStatusUpdateEvent_v0_3; import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Consumer; import java.util.logging.Logger; @@ -17,7 +18,7 @@ public class SSEEventListener_v0_3 { private static final Logger LOGGER = Logger.getLogger(SSEEventListener_v0_3.class.getName()); private final Consumer eventHandler; private final Consumer errorHandler; - private volatile boolean completed = false; + private final AtomicBoolean terminalSignaled = new AtomicBoolean(false); public SSEEventListener_v0_3(Consumer eventHandler, Consumer errorHandler) { @@ -34,44 +35,37 @@ public void onMessage(String message, Future completableFuture) { LOGGER.warning("Failed to process JSON message: " + message); } catch (IllegalArgumentException e) { LOGGER.warning("Invalid message format: " + message); - if (errorHandler != null) { - errorHandler.accept(e); - } + signalTerminal(e); completableFuture.cancel(true); // close SSE channel } } public void onError(Throwable throwable, Future future) { - if (errorHandler != null) { - errorHandler.accept(throwable); - } + signalTerminal(throwable); future.cancel(true); // close SSE channel } - public void onComplete() { - // Idempotent: only signal completion once, even if called multiple times - if (completed) { - LOGGER.fine("SSEEventListener.onComplete() called again - ignoring (already completed)"); + private void signalTerminal(Throwable error) { + if (!terminalSignaled.compareAndSet(false, true)) { + LOGGER.fine("Terminal callback already delivered, ignoring subsequent signal"); return; } - completed = true; - - // Signal normal stream completion (null error means successful completion) - LOGGER.fine("SSEEventListener.onComplete() called - signaling successful stream completion"); if (errorHandler != null) { - LOGGER.fine("Calling errorHandler.accept(null) to signal successful completion"); - errorHandler.accept(null); - } else { - LOGGER.warning("errorHandler is null, cannot signal completion"); + errorHandler.accept(error); + } else if (error != null) { + LOGGER.warning("errorHandler is null, cannot report terminal error"); } } + public void onComplete() { + LOGGER.fine("SSEEventListener.onComplete() called - signaling successful stream completion"); + signalTerminal(null); + } + private void handleMessage(JsonObject jsonObject, Future future) throws JsonProcessingException_v0_3 { if (jsonObject.has("error")) { JSONRPCError_v0_3 error = JsonUtil_v0_3.fromJson(jsonObject.get("error").toString(), JSONRPCError_v0_3.class); - if (errorHandler != null) { - errorHandler.accept(error); - } + signalTerminal(error); } else if (jsonObject.has("result")) { // result can be a Task, Message, TaskStatusUpdateEvent, or TaskArtifactUpdateEvent String resultJson = jsonObject.get("result").toString(); diff --git a/compat-0.3/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3_Test.java b/compat-0.3/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3_Test.java index 8db887db5..c3a248179 100644 --- a/compat-0.3/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3_Test.java +++ b/compat-0.3/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3_Test.java @@ -4,12 +4,18 @@ import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertTrue; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; @@ -232,6 +238,89 @@ public void testOnEventWithFinalTaskStatusUpdateEventEventCancels() throws Excep } + + @Test + public void testOnCompleteThenOnErrorDeliversSingleTerminalCallback() { + AtomicInteger terminalCount = new AtomicInteger(0); + AtomicReference lastArg = new AtomicReference<>(); + SSEEventListener_v0_3 listener = new SSEEventListener_v0_3( + event -> {}, + error -> { terminalCount.incrementAndGet(); lastArg.set(error); }); + + listener.onComplete(); + listener.onError(new RuntimeException("late error"), new CancelCapturingFuture()); + + assertEquals(1, terminalCount.get()); + assertNull(lastArg.get()); + } + + @Test + public void testOnErrorThenOnCompleteDeliversSingleTerminalCallback() { + AtomicInteger terminalCount = new AtomicInteger(0); + AtomicReference lastArg = new AtomicReference<>(); + SSEEventListener_v0_3 listener = new SSEEventListener_v0_3( + event -> {}, + error -> { terminalCount.incrementAndGet(); lastArg.set(error); }); + + RuntimeException boom = new RuntimeException("first error"); + listener.onError(boom, new CancelCapturingFuture()); + listener.onComplete(); + + assertEquals(1, terminalCount.get()); + assertSame(boom, lastArg.get()); + } + + @Test + public void testRepeatedOnCompleteDeliversSingleTerminalCallback() { + AtomicInteger terminalCount = new AtomicInteger(0); + AtomicReference lastArg = new AtomicReference<>(); + SSEEventListener_v0_3 listener = new SSEEventListener_v0_3( + event -> {}, + error -> { terminalCount.incrementAndGet(); lastArg.set(error); }); + + listener.onComplete(); + listener.onComplete(); + listener.onComplete(); + + assertEquals(1, terminalCount.get()); + } + + @Test + public void testConcurrentTerminalSignalsDeliverExactlyOneCallback() throws Exception { + AtomicInteger terminalCount = new AtomicInteger(0); + AtomicReference lastArg = new AtomicReference<>(); + SSEEventListener_v0_3 listener = new SSEEventListener_v0_3( + event -> {}, + error -> { terminalCount.incrementAndGet(); lastArg.set(error); }); + + int threads = 32; + ExecutorService pool = Executors.newFixedThreadPool(threads); + CountDownLatch start = new CountDownLatch(1); + CountDownLatch done = new CountDownLatch(threads); + for (int i = 0; i < threads; i++) { + final boolean complete = (i % 2 == 0); + pool.submit(() -> { + try { + start.await(); + if (complete) { + listener.onComplete(); + } else { + listener.onError(new RuntimeException("race"), new CancelCapturingFuture()); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + done.countDown(); + } + }); + } + start.countDown(); + assertTrue(done.await(10, TimeUnit.SECONDS)); + pool.shutdownNow(); + + assertEquals(1, terminalCount.get()); + } + private static class CancelCapturingFuture implements Future { private boolean cancelHandlerCalled; From 29e414a61bafcb25027cab0ee1c57022c6b170de Mon Sep 17 00:00:00 2001 From: Kabir Khan Date: Thu, 24 Sep 2026 17:37:03 +0100 Subject: [PATCH 2/3] fix(client): signal completion for final SSE events Signal normal completion before auto-closing streams for final events so cancellation cannot report a contradictory terminal outcome. Apply this to the shared listener and the v0.3 JSON-RPC listener. This fixes #1170 --- .../sdk/client/transport/spi/sse/AbstractSSEEventListener.java | 1 + .../client/transport/jsonrpc/sse/SSEEventListener_v0_3.java | 1 + 2 files changed, 2 insertions(+) diff --git a/client/transport/spi/src/main/java/org/a2aproject/sdk/client/transport/spi/sse/AbstractSSEEventListener.java b/client/transport/spi/src/main/java/org/a2aproject/sdk/client/transport/spi/sse/AbstractSSEEventListener.java index 29b30d000..2078bc7ec 100644 --- a/client/transport/spi/src/main/java/org/a2aproject/sdk/client/transport/spi/sse/AbstractSSEEventListener.java +++ b/client/transport/spi/src/main/java/org/a2aproject/sdk/client/transport/spi/sse/AbstractSSEEventListener.java @@ -118,6 +118,7 @@ protected void handleEvent(StreamingEventKind event, @Nullable Future futu // This covers late subscriptions to completed tasks and ensures no connection leaks if (shouldAutoClose(event) && future != null) { LOGGER.fine("Auto-closing SSE connection for final event: " + event.getClass().getSimpleName()); + signalTerminal(null); future.cancel(true); // close SSE channel } } diff --git a/compat-0.3/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3.java b/compat-0.3/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3.java index 22b18a4be..991769a2d 100644 --- a/compat-0.3/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3.java +++ b/compat-0.3/client/transport/jsonrpc/src/main/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3.java @@ -72,6 +72,7 @@ private void handleMessage(JsonObject jsonObject, Future future) throws Js StreamingEventKind_v0_3 event = JsonUtil_v0_3.fromJson(resultJson, StreamingEventKind_v0_3.class); eventHandler.accept(event); if (event instanceof TaskStatusUpdateEvent_v0_3 tsue && tsue.isFinal()) { + signalTerminal(null); future.cancel(true); // close SSE channel } } else { From 13e4f409b4270ca59c881df64cf9b2e32bc6b586 Mon Sep 17 00:00:00 2001 From: Kabir Khan Date: Thu, 24 Sep 2026 17:43:37 +0100 Subject: [PATCH 3/3] test(client): cover final SSE cancellation callback Verify that an onError callback caused by closing a stream after a final event does not replace the normal completion outcome in native and v0.3 JSON-RPC listeners. This fixes #1170 --- .../jsonrpc/sse/SSEEventListenerTest.java | 24 ++++++++++++++++++- .../sse/SSEEventListener_v0_3_Test.java | 24 ++++++++++++++++++- 2 files changed, 46 insertions(+), 2 deletions(-) diff --git a/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListenerTest.java b/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListenerTest.java index a3424d99c..3efc51fcd 100644 --- a/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListenerTest.java +++ b/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/client/transport/jsonrpc/sse/SSEEventListenerTest.java @@ -354,6 +354,28 @@ public void testFinalEventThenOnCompleteDeliversSingleTerminalCallback() throws assertEquals(1, terminalCount.get()); } + @Test + public void testFinalEventThenOnErrorDeliversNormalCompletionOnly() { + AtomicInteger terminalCount = new AtomicInteger(0); + AtomicReference terminalError = new AtomicReference<>(); + SSEEventListener listener = new SSEEventListener( + event -> {}, + error -> { + terminalCount.incrementAndGet(); + terminalError.set(error); + }); + + String eventData = JsonStreamingMessages.STREAMING_STATUS_UPDATE_EVENT_FINAL.substring( + JsonStreamingMessages.STREAMING_STATUS_UPDATE_EVENT_FINAL.indexOf("{")); + CancelCapturingFuture future = new CancelCapturingFuture(); + listener.onMessage(new ServerSentEvent(eventData), future); + listener.onError(new RuntimeException("cancelled after final event"), future); + + assertTrue(future.cancelHandlerCalled); + assertEquals(1, terminalCount.get()); + assertNull(terminalError.get()); + } + private static class CancelCapturingFuture implements Future { private boolean cancelHandlerCalled; @@ -386,4 +408,4 @@ public Void get(long timeout, TimeUnit unit) throws InterruptedException, Execut return null; } } -} \ No newline at end of file +} diff --git a/compat-0.3/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3_Test.java b/compat-0.3/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3_Test.java index c3a248179..0f22828a5 100644 --- a/compat-0.3/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3_Test.java +++ b/compat-0.3/client/transport/jsonrpc/src/test/java/org/a2aproject/sdk/compat03/client/transport/jsonrpc/sse/SSEEventListener_v0_3_Test.java @@ -321,6 +321,28 @@ public void testConcurrentTerminalSignalsDeliverExactlyOneCallback() throws Exce assertEquals(1, terminalCount.get()); } + @Test + public void testFinalEventThenOnErrorDeliversNormalCompletionOnly() { + AtomicInteger terminalCount = new AtomicInteger(0); + AtomicReference terminalError = new AtomicReference<>(); + SSEEventListener_v0_3 listener = new SSEEventListener_v0_3( + event -> {}, + error -> { + terminalCount.incrementAndGet(); + terminalError.set(error); + }); + + String eventData = JsonStreamingMessages_v0_3.STREAMING_STATUS_UPDATE_EVENT_FINAL.substring( + JsonStreamingMessages_v0_3.STREAMING_STATUS_UPDATE_EVENT_FINAL.indexOf("{")); + CancelCapturingFuture future = new CancelCapturingFuture(); + listener.onMessage(eventData, future); + listener.onError(new RuntimeException("cancelled after final event"), future); + + assertTrue(future.cancelHandlerCalled); + assertEquals(1, terminalCount.get()); + assertNull(terminalError.get()); + } + private static class CancelCapturingFuture implements Future { private boolean cancelHandlerCalled; @@ -353,4 +375,4 @@ public Void get(long timeout, TimeUnit unit) throws InterruptedException, Execut return null; } } -} \ No newline at end of file +}