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..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 @@ -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; @@ -201,6 +229,8 @@ public Flow.Publisher consumeAll() { 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); @@ -235,10 +265,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..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 @@ -236,6 +236,50 @@ public void testConsumeMessageEvents() throws Exception { assertSame(message, receivedEvents.get(0)); } + @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 public void testConsumeTaskInputRequired() { // Per A2A Protocol Specification 3.1.6 (SubscribeToTask): 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..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 @@ -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: + // 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(). + assertEquals(0, mainQueue.size(), "Semaphore permit must be released on submit failure"); + } }