From 7eb31d50d686ca935d78c5607c9eafb082324729 Mon Sep 17 00:00:00 2001 From: meraklbz Date: Mon, 10 Aug 2026 20:05:05 +0800 Subject: [PATCH 1/8] fix: harden event consumer and queue --- .../sdk/server/events/EventConsumer.java | 50 ++++++++++++++--- .../sdk/server/events/EventQueue.java | 20 +++++-- .../sdk/server/events/EventConsumerTest.java | 56 ++++++++++++++++++- .../sdk/server/events/EventQueueTest.java | 41 ++++++++++++++ 4 files changed, 153 insertions(+), 14 deletions(-) diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java b/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java index 18da23b22..80f739034 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java @@ -47,7 +47,35 @@ public class EventConsumer { // the write callback, causing onComplete to fire before the HTTP response.write() // callback confirms the data was sent. This sleep ensures the write callback fires // first, so response.end() is only called after the data is safely in flight. - private static final int BUFFER_FLUSH_DELAY_MS = 150; + // + // The delay only applies once per stream, when the final event is sent. It is + // configurable via the system property {@value #BUFFER_FLUSH_DELAY_MS_PROPERTY} + // (milliseconds, default 150, minimum 0) so operators can trade flush reliability + // against stream-termination latency for their transport. + private static final int DEFAULT_BUFFER_FLUSH_DELAY_MS = 150; + private static final String BUFFER_FLUSH_DELAY_MS_PROPERTY = "a2a.eventconsumer.bufferFlushDelayMs"; + + /** + * Returns the configured buffer-flush delay in milliseconds. + * + *

Reads the {@value #BUFFER_FLUSH_DELAY_MS_PROPERTY} system property; values that + * are absent, non-numeric, or negative fall back to the default.

+ * + * @return the delay in milliseconds (never negative) + */ + static int bufferFlushDelayMs() { + String configured = System.getProperty(BUFFER_FLUSH_DELAY_MS_PROPERTY); + if (configured == null) { + return DEFAULT_BUFFER_FLUSH_DELAY_MS; + } + try { + return Math.max(0, Integer.parseInt(configured.trim())); + } catch (NumberFormatException e) { + LOGGER.warn("Invalid {} value '{}', falling back to default {}", + BUFFER_FLUSH_DELAY_MS_PROPERTY, configured, DEFAULT_BUFFER_FLUSH_DELAY_MS); + return DEFAULT_BUFFER_FLUSH_DELAY_MS; + } + } public EventConsumer(EventQueue queue, Executor executor) { this.queue = queue; @@ -200,8 +228,6 @@ public Flow.Publisher consumeAll() { boolean isFinalEvent = false; if (event instanceof TaskStatusUpdateEvent tue && tue.isFinal()) { isFinalEvent = true; - } else if (event instanceof Message) { - isFinalEvent = true; } else if (event instanceof Task task) { isFinalEvent = isStreamTerminatingTask(task); } else if (event instanceof QueueClosedEvent) { @@ -215,6 +241,13 @@ public Flow.Publisher consumeAll() { LOGGER.debug("Received A2AError event, treating as final event"); isFinalEvent = true; } + // NOTE: A plain Message event is intentionally NOT stream-terminating. + // Per the A2A protocol the stream MUST terminate only when the task reaches + // a terminal state (completed, failed, canceled, rejected); an intermediate + // message emitted before the agent finishes (or a message-only response that + // still has follow-up events) must not close the stream early. The stream is + // closed by a final status update/task, a QueueClosedEvent, an A2AError, or + // the agent-completed grace period when no final event arrives. // Only send event if it's not a QueueClosedEvent // QueueClosedEvent is an internal coordination event used for replication @@ -235,10 +268,13 @@ public Flow.Publisher consumeAll() { // of the write callback, causing response.end() to race with a pending // response.write(). This delay ensures the write callback runs first. if (isFinalSent) { - try { - Thread.sleep(BUFFER_FLUSH_DELAY_MS); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); + int flushDelayMs = bufferFlushDelayMs(); + if (flushDelayMs > 0) { + try { + Thread.sleep(flushDelayMs); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } } } break; diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/events/EventQueue.java b/server-common/src/main/java/org/a2aproject/sdk/server/events/EventQueue.java index 4fb8ad4f6..3c24a103d 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/events/EventQueue.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/events/EventQueue.java @@ -532,9 +532,17 @@ public void enqueueItem(EventQueueItem item) { // Submit to MainEventBus for centralized persistence + distribution // MainEventBus is guaranteed non-null by constructor requirement // Note: Replication now happens in MainEventBusProcessor AFTER persistence - - // Submit event to MainEventBus with our taskId - mainEventBus.submit(taskId, this, item); + try { + // Submit event to MainEventBus with our taskId + mainEventBus.submit(taskId, this, item); + } catch (RuntimeException e) { + // The event never reached MainEventBusProcessor, so it will never call + // releaseSemaphore() (see MainEventBusProcessor.processEvent finally block). + // Release the permit here to avoid leaking it and eventually blocking + // all event processing for this task. + semaphore.release(); + throw e; + } } /** @@ -766,12 +774,16 @@ String getTaskId() { static class ChildQueue extends EventQueue { private final MainQueue parent; - private final BlockingQueue queue = new LinkedBlockingDeque<>(); + private final BlockingQueue queue; private volatile boolean immediateClose = false; private volatile boolean awaitingFinalEvent = false; public ChildQueue(MainQueue parent) { this.parent = parent; + // Bound the child queue with the same capacity as its parent so that a slow + // subscriber cannot grow the deque unboundedly. internalEnqueueItem() detects + // a full queue and closes the child immediately (see offer() below). + this.queue = new LinkedBlockingDeque<>(parent.getQueueSize()); } @Override diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java index 41389c3ec..6b2b561d7 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java @@ -221,9 +221,14 @@ public void testConsumeMessageEvents() throws Exception { List events = List.of(message, message2); for (Event event : events) { - eventQueue.enqueueEvent(event); + waitForEventProcessing(() -> eventQueue.enqueueEvent(event)); } + // A plain Message is no longer treated as a stream-terminating final event + // (BUG-26): the stream must stay open for subsequent events. Close the queue + // explicitly so the polling loop terminates after draining the messages. + eventQueue.close(); + Flow.Publisher publisher = eventConsumer.consumeAll(); final List receivedEvents = new ArrayList<>(); final AtomicReference error = new AtomicReference<>(); @@ -231,9 +236,54 @@ public void testConsumeMessageEvents() throws Exception { publisher.subscribe(getSubscriber(receivedEvents, error)); assertNull(error.get()); - // The stream is closed after the first Message - assertEquals(1, receivedEvents.size()); + // Both messages should be delivered - the stream is no longer closed by the first Message + assertEquals(2, receivedEvents.size()); assertSame(message, receivedEvents.get(0)); + assertSame(message2, receivedEvents.get(1)); + } + + @Test + public void testBufferFlushDelayMsDefaultsTo150() { + String original = System.getProperty("a2a.eventconsumer.bufferFlushDelayMs"); + try { + System.clearProperty("a2a.eventconsumer.bufferFlushDelayMs"); + assertEquals(150, EventConsumer.bufferFlushDelayMs()); + } finally { + restoreProperty("a2a.eventconsumer.bufferFlushDelayMs", original); + } + } + + @Test + public void testBufferFlushDelayMsReadsConfiguredValue() { + String original = System.getProperty("a2a.eventconsumer.bufferFlushDelayMs"); + try { + System.setProperty("a2a.eventconsumer.bufferFlushDelayMs", "20"); + assertEquals(20, EventConsumer.bufferFlushDelayMs()); + } finally { + restoreProperty("a2a.eventconsumer.bufferFlushDelayMs", original); + } + } + + @Test + public void testBufferFlushDelayMsRejectsInvalidValues() { + String original = System.getProperty("a2a.eventconsumer.bufferFlushDelayMs"); + try { + System.setProperty("a2a.eventconsumer.bufferFlushDelayMs", "not-a-number"); + assertEquals(150, EventConsumer.bufferFlushDelayMs()); + // Negative values are clamped to 0 (disabled) + System.setProperty("a2a.eventconsumer.bufferFlushDelayMs", "-5"); + assertEquals(0, EventConsumer.bufferFlushDelayMs()); + } finally { + restoreProperty("a2a.eventconsumer.bufferFlushDelayMs", original); + } + } + + private static void restoreProperty(String key, String original) { + if (original == null) { + System.clearProperty(key); + } else { + System.setProperty(key, original); + } } @Test diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/events/EventQueueTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/events/EventQueueTest.java index 2693b008a..8730468b5 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/events/EventQueueTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/events/EventQueueTest.java @@ -702,4 +702,45 @@ public void onTaskFinalized(String taskId) { assertFalse(onEventCalled.get(), "onEvent should not be called when there are zero subscribers"); } + + @Test + public void testChildQueueIsBoundedByParentQueueSize() throws Exception { + int customSize = 5; + EventQueue mainQueue = EventQueueUtil.getEventQueueBuilder(mainEventBus) + .queueSize(customSize) + .build(); + EventQueue childQueue = mainQueue.tap(); + + java.lang.reflect.Field queueField = EventQueue.ChildQueue.class.getDeclaredField("queue"); + queueField.setAccessible(true); + java.util.concurrent.BlockingQueue childDeque = + (java.util.concurrent.BlockingQueue) queueField.get(childQueue); + + // The child queue must be bounded by the parent's configured capacity (BUG-28): + // an unbounded deque would let a slow subscriber grow memory without limit. + assertEquals(customSize, childDeque.remainingCapacity()); + } + + @Test + public void testSemaphorePermitReleasedWhenSubmitFails() { + MainEventBus failingBus = org.mockito.Mockito.mock(MainEventBus.class); + org.mockito.Mockito.doThrow(new RuntimeException("submit failed")) + .when(failingBus).submit(org.mockito.ArgumentMatchers.anyString(), + org.mockito.ArgumentMatchers.any(), + org.mockito.ArgumentMatchers.any()); + + EventQueue mainQueue = EventQueueUtil.getEventQueueBuilder(failingBus) + .taskId(TASK_ID) + .queueSize(2) + .build(); + assertEquals(0, mainQueue.size(), "No permits should be in use before enqueue"); + + assertThrows(RuntimeException.class, + () -> mainQueue.enqueueEvent(fromJson(MINIMAL_TASK, Task.class))); + + // If the permit leaked, size() would be 1 (one permit held forever). It must be + // released when submit() fails because MainEventBusProcessor never saw the event + // and therefore never called releaseSemaphore() (BUG-29). + assertEquals(0, mainQueue.size(), "Semaphore permit must be released on submit failure"); + } } From cd1d5d711e23483b1c6d69bde101b9829ec35019 Mon Sep 17 00:00:00 2001 From: meraklbz Date: Tue, 11 Aug 2026 03:35:10 +0800 Subject: [PATCH 2/8] ci: retrigger upstream checks From 2178a48f15f06675510a599ed7b12d46ee9a9713 Mon Sep 17 00:00:00 2001 From: meraklbz Date: Tue, 11 Aug 2026 04:09:24 +0800 Subject: [PATCH 3/8] ci: retrigger checks (stuck run workaround) From 0be13c607ba9e49cee05587b2334b19489488dad Mon Sep 17 00:00:00 2001 From: meraklbz Date: Tue, 11 Aug 2026 04:32:37 +0800 Subject: [PATCH 4/8] ci: retrigger stuck build matrix From 00a72cc6fd8792d9a2697c174d45f35be7cee0b3 Mon Sep 17 00:00:00 2001 From: meraklbz Date: Tue, 11 Aug 2026 07:51:42 +0800 Subject: [PATCH 5/8] ci: retrigger build matrix From 2c42fa7ce7992249ead742ffdda4a500c81c2d90 Mon Sep 17 00:00:00 2001 From: meraklbz Date: Tue, 11 Aug 2026 21:17:15 +0800 Subject: [PATCH 6/8] =?UTF-8?q?fix:=20keep=20Message=20as=20a=20stream-ter?= =?UTF-8?q?minating=20event=20per=20A2A=20spec=20=C2=A73.1.2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Revert the change that treated Message events as non-terminating. Per the A2A spec (Send Streaming Message), a Message is the complete response and the stream must close after delivering it. The EventConsumer serves both message/stream and task subscriptions, and the spec-conformant behavior for message/stream takes precedence. The BUG-26-motivated test asserting the opposite is removed. Keep the other three changes: configurable buffer-flush delay, bounded ChildQueue, and the semaphore permit leak fix. --- .../sdk/server/events/EventConsumer.java | 4 +++ .../sdk/server/events/EventConsumerTest.java | 29 ------------------- 2 files changed, 4 insertions(+), 29 deletions(-) diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java b/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java index 80f739034..4c1ffa7f9 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java @@ -228,6 +228,10 @@ public Flow.Publisher consumeAll() { boolean isFinalEvent = false; if (event instanceof TaskStatusUpdateEvent tue && tue.isFinal()) { isFinalEvent = true; + } else if (event instanceof Message) { + // Per A2A spec §3.1.2 (Send Streaming Message): a Message is the + // complete response — the stream must close after delivering it. + isFinalEvent = true; } else if (event instanceof Task task) { isFinalEvent = isStreamTerminatingTask(task); } else if (event instanceof QueueClosedEvent) { diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java index 6b2b561d7..c26952980 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java @@ -213,35 +213,6 @@ public void testConsumeUntilMessage() throws Exception { } } - @Test - public void testConsumeMessageEvents() throws Exception { - Message message = fromJson(MESSAGE_PAYLOAD, Message.class); - Message message2 = Message.builder(message).build(); - - List events = List.of(message, message2); - - for (Event event : events) { - waitForEventProcessing(() -> eventQueue.enqueueEvent(event)); - } - - // A plain Message is no longer treated as a stream-terminating final event - // (BUG-26): the stream must stay open for subsequent events. Close the queue - // explicitly so the polling loop terminates after draining the messages. - eventQueue.close(); - - Flow.Publisher publisher = eventConsumer.consumeAll(); - final List receivedEvents = new ArrayList<>(); - final AtomicReference error = new AtomicReference<>(); - - publisher.subscribe(getSubscriber(receivedEvents, error)); - - assertNull(error.get()); - // Both messages should be delivered - the stream is no longer closed by the first Message - assertEquals(2, receivedEvents.size()); - assertSame(message, receivedEvents.get(0)); - assertSame(message2, receivedEvents.get(1)); - } - @Test public void testBufferFlushDelayMsDefaultsTo150() { String original = System.getProperty("a2a.eventconsumer.bufferFlushDelayMs"); From ef85f71922ce3bbbd7dce65fe9902fb92ea0241b Mon Sep 17 00:00:00 2001 From: meraklbz Date: Tue, 11 Aug 2026 21:38:05 +0800 Subject: [PATCH 7/8] chore: remove internal tracking ids from comments --- .../java/org/a2aproject/sdk/server/events/EventQueueTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/events/EventQueueTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/events/EventQueueTest.java index 8730468b5..6df08a7e1 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/events/EventQueueTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/events/EventQueueTest.java @@ -716,7 +716,7 @@ public void testChildQueueIsBoundedByParentQueueSize() throws Exception { java.util.concurrent.BlockingQueue childDeque = (java.util.concurrent.BlockingQueue) queueField.get(childQueue); - // The child queue must be bounded by the parent's configured capacity (BUG-28): + // The child queue must be bounded by the parent's configured capacity: // an unbounded deque would let a slow subscriber grow memory without limit. assertEquals(customSize, childDeque.remainingCapacity()); } @@ -740,7 +740,7 @@ public void testSemaphorePermitReleasedWhenSubmitFails() { // If the permit leaked, size() would be 1 (one permit held forever). It must be // released when submit() fails because MainEventBusProcessor never saw the event - // and therefore never called releaseSemaphore() (BUG-29). + // and therefore never called releaseSemaphore(). assertEquals(0, mainQueue.size(), "Semaphore permit must be released on submit failure"); } } From 841970c4cdd59dfd5d4ea91af3f77a29a3f98ce9 Mon Sep 17 00:00:00 2001 From: Kabir Khan Date: Tue, 11 Aug 2026 15:52:17 +0100 Subject: [PATCH 8/8] fix: remove stale comment and restore deleted test from revert The revert of the "Message is not stream-terminating" change left behind a contradictory NOTE comment and didn't restore testConsumeMessageEvents. Co-Authored-By: Claude Opus 4.6 (1M context) --- .../sdk/server/events/EventConsumer.java | 7 ------ .../sdk/server/events/EventConsumerTest.java | 23 +++++++++++++++++++ 2 files changed, 23 insertions(+), 7 deletions(-) diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java b/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java index 4c1ffa7f9..8719ff611 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/events/EventConsumer.java @@ -245,13 +245,6 @@ public Flow.Publisher consumeAll() { LOGGER.debug("Received A2AError event, treating as final event"); isFinalEvent = true; } - // NOTE: A plain Message event is intentionally NOT stream-terminating. - // Per the A2A protocol the stream MUST terminate only when the task reaches - // a terminal state (completed, failed, canceled, rejected); an intermediate - // message emitted before the agent finishes (or a message-only response that - // still has follow-up events) must not close the stream early. The stream is - // closed by a final status update/task, a QueueClosedEvent, an A2AError, or - // the agent-completed grace period when no final event arrives. // Only send event if it's not a QueueClosedEvent // QueueClosedEvent is an internal coordination event used for replication diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java index c26952980..54882a078 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/events/EventConsumerTest.java @@ -213,6 +213,29 @@ public void testConsumeUntilMessage() throws Exception { } } + @Test + public void testConsumeMessageEvents() throws Exception { + Message message = fromJson(MESSAGE_PAYLOAD, Message.class); + Message message2 = Message.builder(message).build(); + + List events = List.of(message, message2); + + for (Event event : events) { + eventQueue.enqueueEvent(event); + } + + Flow.Publisher publisher = eventConsumer.consumeAll(); + final List receivedEvents = new ArrayList<>(); + final AtomicReference error = new AtomicReference<>(); + + publisher.subscribe(getSubscriber(receivedEvents, error)); + + assertNull(error.get()); + // The stream is closed after the first Message + assertEquals(1, receivedEvents.size()); + assertSame(message, receivedEvents.get(0)); + } + @Test public void testBufferFlushDelayMsDefaultsTo150() { String original = System.getProperty("a2a.eventconsumer.bufferFlushDelayMs");