From 44d0f6ec29e5ae198583f6fd49c631dc712f11e0 Mon Sep 17 00:00:00 2001 From: King Star Date: Thu, 23 Jul 2026 08:44:49 +0800 Subject: [PATCH 1/4] fix(server): reconcile blocking result with TaskStore Fixes a2aproject/A2A#2049 --- .../sdk/server/tasks/ResultAggregator.java | 30 +++++++- .../server/tasks/ResultAggregatorTest.java | 75 +++++++++++++++++++ 2 files changed, 103 insertions(+), 2 deletions(-) diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java b/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java index eddc57bbe..82d9f8e5d 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java @@ -8,6 +8,7 @@ import java.util.concurrent.CompletionException; import java.util.concurrent.Executor; import java.util.concurrent.Flow; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; @@ -28,6 +29,8 @@ public class ResultAggregator { private static final Logger LOGGER = LoggerFactory.getLogger(ResultAggregator.class); + private static final long TASK_STORE_RECONCILIATION_TIMEOUT_NANOS = TimeUnit.SECONDS.toNanos(1); + private static final long TASK_STORE_RECONCILIATION_POLL_MILLIS = 10; private final TaskManager taskManager; private final Executor executor; @@ -235,7 +238,7 @@ else if (blocking) { Utils.rethrow(error); } - // Return Message if captured, otherwise Task if captured, otherwise fetch from TaskStore + // Return Message if captured, otherwise Task if captured, otherwise reconcile with TaskStore. EventKind eventKind = message.get(); if (eventKind == null) { eventKind = capturedTask.get(); @@ -244,7 +247,7 @@ else if (blocking) { } } if (eventKind == null) { - eventKind = taskManager.getTask(); + eventKind = reconcileTaskStore(blocking); if (LOGGER.isDebugEnabled() && eventKind instanceof Task t) { LOGGER.debug("Returning task from TaskStore: id={}, state={}", t.id(), t.status().state()); } @@ -259,6 +262,29 @@ else if (blocking) { consumptionCompletionFuture); } + private @Nullable Task reconcileTaskStore(boolean blocking) throws A2AError { + Task task = taskManager.getTask(); + if (task != null || !blocking) { + return task; + } + + long deadline = System.nanoTime() + TASK_STORE_RECONCILIATION_TIMEOUT_NANOS; + while (System.nanoTime() < deadline) { + try { + Thread.sleep(TASK_STORE_RECONCILIATION_POLL_MILLIS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new InternalError("Interrupted while reconciling TaskStore for " + taskManager.getTaskId()); + } + + task = taskManager.getTask(); + if (task != null) { + return task; + } + } + return null; + } + private String taskIdForLogging() { Task task = taskManager.getTask(); return task != null ? task.id() : "unknown"; diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java index 7b206e01a..329cdb617 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java @@ -3,8 +3,10 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.atMost; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; @@ -15,7 +17,9 @@ import java.util.concurrent.Executor; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import org.a2aproject.sdk.server.events.EnhancedRunnable; import org.a2aproject.sdk.server.events.EventConsumer; import org.a2aproject.sdk.server.events.EventQueue; import org.a2aproject.sdk.server.events.EventQueueUtil; @@ -24,6 +28,7 @@ import org.a2aproject.sdk.server.events.MainEventBusProcessor; import org.a2aproject.sdk.spec.Event; import org.a2aproject.sdk.spec.EventKind; +import org.a2aproject.sdk.spec.InternalError; import org.a2aproject.sdk.spec.Message; import org.a2aproject.sdk.spec.Task; import org.a2aproject.sdk.spec.TaskState; @@ -281,6 +286,76 @@ void testConsumeAndBreakNonBlocking() throws Exception { EventQueueUtil.stop(processor); } + @Test + void testBlockingReturnsTaskWhenStoreBecomesVisibleAfterEmptyCapture() throws Exception { + String taskId = "task-store-reconciliation"; + Task completedTask = createSampleTask(taskId, TaskState.TASK_STATE_COMPLETED, "ctx1"); + AtomicBoolean taskVisible = new AtomicBoolean(false); + TaskStore taskStore = mock(TaskStore.class); + when(taskStore.get(taskId)).thenAnswer(invocation -> + taskVisible.compareAndSet(false, true) ? null : completedTask); + + TaskManager taskManager = new TaskManager(taskId, "ctx1", taskStore, null); + ResultAggregator blockingAggregator = + new ResultAggregator(taskManager, null, testExecutor, testExecutor); + + MainEventBus mainEventBus = new MainEventBus(); + InMemoryQueueManager queueManager = + new InMemoryQueueManager(new MockTaskStateProvider(), mainEventBus); + EventQueue queue = queueManager.getEventQueueBuilder(taskId).build().tap(); + EventConsumer eventConsumer = new EventConsumer(queue, testExecutor); + + EnhancedRunnable completedAgent = new EnhancedRunnable() { + @Override + public void run() { + } + }; + eventConsumer.createAgentRunnableDoneCallback().done(completedAgent); + + ResultAggregator.EventTypeAndInterrupt result = + blockingAggregator.consumeAndBreakOnInterrupt(eventConsumer, true); + + assertEquals(completedTask, result.eventType()); + assertFalse(result.interrupted()); + } + + @Test + void testBlockingStillRaisesMissingTaskErrorWhenStoreRemainsEmpty() { + String taskId = "missing-task"; + TaskStore taskStore = mock(TaskStore.class); + TaskManager taskManager = new TaskManager(taskId, "ctx1", taskStore, null); + ResultAggregator blockingAggregator = + new ResultAggregator(taskManager, null, testExecutor, testExecutor); + + InternalError error = assertThrows(InternalError.class, () -> + blockingAggregator.consumeAndBreakOnInterrupt(createClosedEventConsumer(taskId), true)); + + assertEquals("Could not find a Task/Message for " + taskId, error.getMessage()); + } + + @Test + void testNonBlockingDoesNotPollTaskStoreWhenCaptureIsEmpty() { + String taskId = "non-blocking-missing-task"; + TaskStore taskStore = mock(TaskStore.class); + TaskManager taskManager = new TaskManager(taskId, "ctx1", taskStore, null); + ResultAggregator nonBlockingAggregator = + new ResultAggregator(taskManager, null, testExecutor, testExecutor); + + assertThrows(InternalError.class, () -> + nonBlockingAggregator.consumeAndBreakOnInterrupt(createClosedEventConsumer(taskId), false)); + + verify(taskStore, times(1)).get(taskId); + } + + private EventConsumer createClosedEventConsumer(String taskId) { + MainEventBus mainEventBus = new MainEventBus(); + InMemoryQueueManager queueManager = + new InMemoryQueueManager(new MockTaskStateProvider(), mainEventBus); + EventQueue queue = queueManager.getEventQueueBuilder(taskId).build().tap(); + queue.close(); + return new EventConsumer(queue, testExecutor); + } + // AUTH_REQUIRED Tests @Test From b21bef669afc8492c7833946bdf37664e878f616 Mon Sep 17 00:00:00 2001 From: Kabir Khan Date: Thu, 23 Jul 2026 09:41:11 +0100 Subject: [PATCH 2/4] fix(server): make reconciliation timeout configurable, add logging Co-Authored-By: Claude Opus 4.6 (1M context) --- docs/content/configuration.md | 4 +++ integrations/microprofile-config/README.md | 1 + .../DefaultRequestHandler.java | 34 +++++++++++++++---- .../sdk/server/tasks/ResultAggregator.java | 15 +++++--- .../META-INF/a2a-defaults.properties | 4 +++ .../server/tasks/ResultAggregatorTest.java | 16 ++++----- 6 files changed, 56 insertions(+), 18 deletions(-) diff --git a/docs/content/configuration.md b/docs/content/configuration.md index b164bf3cf..70d337bc8 100644 --- a/docs/content/configuration.md +++ b/docs/content/configuration.md @@ -40,6 +40,9 @@ a2a.blocking.agent.timeout.seconds=30 # Timeout for event consumption in blocking calls (default: 5 seconds) a2a.blocking.consumption.timeout.seconds=5 + +# Timeout for TaskStore reconciliation polling in blocking calls (default: 1 second) +a2a.blocking.reconciliation.timeout.seconds=1 ``` ### Tuning Guidelines @@ -48,6 +51,7 @@ a2a.blocking.consumption.timeout.seconds=5 - **Resource Management**: The dedicated executor prevents streaming operations from competing with the ForkJoinPool. - **Concurrency**: In production with high concurrent streaming, increase pool sizes accordingly. - **Agent Timeouts**: LLM-based agents may need longer timeouts (60-120s) compared to simple agents. +- **Reconciliation Timeout**: Increase if blocking calls fail with "Could not find a Task/Message" under heavy load or with slow TaskStore implementations. ## MicroProfile Config Integration diff --git a/integrations/microprofile-config/README.md b/integrations/microprofile-config/README.md index ed9fd74ac..351218641 100644 --- a/integrations/microprofile-config/README.md +++ b/integrations/microprofile-config/README.md @@ -37,6 +37,7 @@ a2a.executor.max-pool-size=100 # Timeout configuration a2a.blocking.agent.timeout.seconds=60 a2a.blocking.consumption.timeout.seconds=10 +a2a.blocking.reconciliation.timeout.seconds=2 ``` **Environment variables:** diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/requesthandlers/DefaultRequestHandler.java b/server-common/src/main/java/org/a2aproject/sdk/server/requesthandlers/DefaultRequestHandler.java index db7478fb0..addac8e69 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/requesthandlers/DefaultRequestHandler.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/requesthandlers/DefaultRequestHandler.java @@ -142,8 +142,9 @@ *
  • Blocking (configuration.blocking=true): Client waits for first event or final task state
  • *
  • Streaming: Client receives events as they arrive via reactive streams
  • *
  • Both modes support fire-and-forget (agent continues after client disconnect)
  • - *
  • Configurable timeouts via {@code a2a.blocking.agent.timeout.seconds} and - * {@code a2a.blocking.consumption.timeout.seconds}
  • + *
  • Configurable timeouts via {@code a2a.blocking.agent.timeout.seconds}, + * {@code a2a.blocking.consumption.timeout.seconds}, and + * {@code a2a.blocking.reconciliation.timeout.seconds}
  • * * *

    CDI Dependencies

    @@ -189,6 +190,7 @@ public class DefaultRequestHandler implements RequestHandler { private static final String A2A_BLOCKING_AGENT_TIMEOUT_SECONDS = "a2a.blocking.agent.timeout.seconds"; private static final String A2A_BLOCKING_CONSUMPTION_TIMEOUT_SECONDS = "a2a.blocking.consumption.timeout.seconds"; + private static final String A2A_BLOCKING_RECONCILIATION_TIMEOUT_SECONDS = "a2a.blocking.reconciliation.timeout.seconds"; @Inject A2AConfigProvider configProvider; @@ -215,6 +217,19 @@ public class DefaultRequestHandler implements RequestHandler { */ int consumptionCompletionTimeoutSeconds; + /** + * Timeout in seconds for TaskStore reconciliation polling in blocking calls. + * When the in-memory event capture is empty, the aggregator polls TaskStore + * with a bounded timeout to handle the race where MainEventBusProcessor + * has not yet persisted the task. + *

    + * Property: {@code a2a.blocking.reconciliation.timeout.seconds}
    + * Default: 1 second
    + * Note: Property override requires a configurable {@link A2AConfigProvider} on the classpath + * (e.g., MicroProfileConfigProvider in reference implementations). + */ + int reconciliationTimeoutSeconds; + // Fields set by constructor injection cannot be final. We need a noargs constructor for // Jakarta compatibility, and it seems that making fields set by constructor injection // final, is not proxyable in all runtimes @@ -276,6 +291,8 @@ void initConfig() { configProvider.getValue(A2A_BLOCKING_AGENT_TIMEOUT_SECONDS)); consumptionCompletionTimeoutSeconds = Integer.parseInt( configProvider.getValue(A2A_BLOCKING_CONSUMPTION_TIMEOUT_SECONDS)); + reconciliationTimeoutSeconds = Integer.parseInt( + configProvider.getValue(A2A_BLOCKING_RECONCILIATION_TIMEOUT_SECONDS)); } @@ -291,6 +308,7 @@ public static DefaultRequestHandler create(AgentExecutor agentExecutor, TaskStor mainEventBusProcessor, executor, eventConsumerExecutor); handler.agentCompletionTimeoutSeconds = 5; handler.consumptionCompletionTimeoutSeconds = 2; + handler.reconciliationTimeoutSeconds = 1; return handler; } @@ -377,7 +395,8 @@ public Task onCancelTask(CancelTaskParams params, ServerCallContext context) thr taskStore, null); - ResultAggregator resultAggregator = new ResultAggregator(taskManager, null, executor, eventConsumerExecutor); + ResultAggregator resultAggregator = new ResultAggregator(taskManager, null, executor, eventConsumerExecutor, + SECONDS.toNanos(reconciliationTimeoutSeconds)); EventQueue queue = queueManager.createOrTap(task.id()); EventConsumer consumer = new EventConsumer(queue, eventConsumerExecutor); @@ -450,7 +469,8 @@ public EventKind onMessageSend(MessageSendParams params, ServerCallContext conte // Create queue with real taskId (no tempId parameter needed) EventQueue queue = queueManager.createOrTap(queueTaskId); final java.util.concurrent.atomic.AtomicReference<@NonNull String> taskId = new java.util.concurrent.atomic.AtomicReference<>(queueTaskId); - ResultAggregator resultAggregator = new ResultAggregator(mss.taskManager, null, executor, eventConsumerExecutor); + ResultAggregator resultAggregator = new ResultAggregator(mss.taskManager, null, executor, eventConsumerExecutor, + SECONDS.toNanos(reconciliationTimeoutSeconds)); // Default to blocking per A2A spec (returnImmediately defaults to false, meaning wait for completion) boolean returnImmediately = params.configuration() != null && Boolean.TRUE.equals(params.configuration().returnImmediately()); @@ -668,7 +688,8 @@ public Flow.Publisher onMessageSendStream( .taskId(taskId.get()).build(), version); } - ResultAggregator resultAggregator = new ResultAggregator(mss.taskManager, null, executor, eventConsumerExecutor); + ResultAggregator resultAggregator = new ResultAggregator(mss.taskManager, null, executor, eventConsumerExecutor, + SECONDS.toNanos(reconciliationTimeoutSeconds)); // Create consumer BEFORE starting agent - callback is registered inside registerAndExecuteAgentAsync EventConsumer consumer = new EventConsumer(queue, eventConsumerExecutor); @@ -857,7 +878,8 @@ public Flow.Publisher onSubscribeToTask(TaskIdParams params, } TaskManager taskManager = new TaskManager(task.id(), task.contextId(), taskStore, null); - ResultAggregator resultAggregator = new ResultAggregator(taskManager, null, executor, eventConsumerExecutor); + ResultAggregator resultAggregator = new ResultAggregator(taskManager, null, executor, eventConsumerExecutor, + SECONDS.toNanos(reconciliationTimeoutSeconds)); EventQueue queue = queueManager.tap(task.id()); LOGGER.debug("onSubscribeToTask - tapped queue: {}", queue != null ? System.identityHashCode(queue) : "null"); diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java b/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java index 82d9f8e5d..32c0b4ad0 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java @@ -29,19 +29,21 @@ public class ResultAggregator { private static final Logger LOGGER = LoggerFactory.getLogger(ResultAggregator.class); - private static final long TASK_STORE_RECONCILIATION_TIMEOUT_NANOS = TimeUnit.SECONDS.toNanos(1); private static final long TASK_STORE_RECONCILIATION_POLL_MILLIS = 10; private final TaskManager taskManager; private final Executor executor; private final Executor eventConsumerExecutor; + private final long reconciliationTimeoutNanos; private volatile @Nullable Message message; - public ResultAggregator(TaskManager taskManager, @Nullable Message message, Executor executor, Executor eventConsumerExecutor) { + public ResultAggregator(TaskManager taskManager, @Nullable Message message, Executor executor, + Executor eventConsumerExecutor, long reconciliationTimeoutNanos) { this.taskManager = taskManager; this.message = message; this.executor = executor; this.eventConsumerExecutor = eventConsumerExecutor; + this.reconciliationTimeoutNanos = reconciliationTimeoutNanos; } public @Nullable EventKind getCurrentResult() { @@ -238,7 +240,8 @@ else if (blocking) { Utils.rethrow(error); } - // Return Message if captured, otherwise Task if captured, otherwise reconcile with TaskStore. + // Return Message if captured, otherwise Task if captured, + // otherwise poll TaskStore with bounded timeout (blocking) or single read (non-blocking). EventKind eventKind = message.get(); if (eventKind == null) { eventKind = capturedTask.get(); @@ -268,7 +271,10 @@ else if (blocking) { return task; } - long deadline = System.nanoTime() + TASK_STORE_RECONCILIATION_TIMEOUT_NANOS; + LOGGER.debug("TaskStore reconciliation: task not found on first read, polling for up to {}ms", + TimeUnit.NANOSECONDS.toMillis(reconciliationTimeoutNanos)); + + long deadline = System.nanoTime() + reconciliationTimeoutNanos; while (System.nanoTime() < deadline) { try { Thread.sleep(TASK_STORE_RECONCILIATION_POLL_MILLIS); @@ -279,6 +285,7 @@ else if (blocking) { task = taskManager.getTask(); if (task != null) { + LOGGER.debug("TaskStore reconciliation: task {} found after polling", task.id()); return task; } } diff --git a/server-common/src/main/resources/META-INF/a2a-defaults.properties b/server-common/src/main/resources/META-INF/a2a-defaults.properties index 719be9e7a..9e0d13480 100644 --- a/server-common/src/main/resources/META-INF/a2a-defaults.properties +++ b/server-common/src/main/resources/META-INF/a2a-defaults.properties @@ -10,6 +10,10 @@ a2a.blocking.agent.timeout.seconds=30 # Ensures TaskStore is fully updated before returning to client a2a.blocking.consumption.timeout.seconds=5 +# Timeout for TaskStore reconciliation polling in blocking calls (seconds) +# When in-memory event capture is empty, polls TaskStore with this bounded timeout +a2a.blocking.reconciliation.timeout.seconds=1 + # AsyncExecutorProducer - Thread pool configuration # Core pool size for async agent execution a2a.executor.core-pool-size=5 diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java index 329cdb617..d3adc986e 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java @@ -62,7 +62,7 @@ public class ResultAggregatorTest { @BeforeEach void setUp() { MockitoAnnotations.openMocks(this); - aggregator = new ResultAggregator(mockTaskManager, null, testExecutor, testExecutor); + aggregator = new ResultAggregator(mockTaskManager, null, testExecutor, testExecutor, TimeUnit.SECONDS.toNanos(1)); } // Helper methods for creating sample data @@ -120,7 +120,7 @@ public void onTaskFinalized(String taskId) { @Test void testConstructorWithMessage() { Message initialMessage = createSampleMessage("initial", "msg1", Message.Role.ROLE_USER); - ResultAggregator aggregatorWithMessage = new ResultAggregator(mockTaskManager, initialMessage, testExecutor, testExecutor); + ResultAggregator aggregatorWithMessage = new ResultAggregator(mockTaskManager, initialMessage, testExecutor, testExecutor, TimeUnit.SECONDS.toNanos(1)); // Test that the message is properly stored by checking getCurrentResult assertEquals(initialMessage, aggregatorWithMessage.getCurrentResult()); @@ -131,7 +131,7 @@ void testConstructorWithMessage() { @Test void testGetCurrentResultWithMessageSet() { Message sampleMessage = createSampleMessage("hola", "msg1", Message.Role.ROLE_USER); - ResultAggregator aggregatorWithMessage = new ResultAggregator(mockTaskManager, sampleMessage, testExecutor, testExecutor); + ResultAggregator aggregatorWithMessage = new ResultAggregator(mockTaskManager, sampleMessage, testExecutor, testExecutor, TimeUnit.SECONDS.toNanos(1)); EventKind result = aggregatorWithMessage.getCurrentResult(); @@ -166,7 +166,7 @@ void testConstructorStoresTaskManagerCorrectly() { @Test void testConstructorWithNullMessage() { - ResultAggregator aggregatorWithNullMessage = new ResultAggregator(mockTaskManager, null, testExecutor, testExecutor); + ResultAggregator aggregatorWithNullMessage = new ResultAggregator(mockTaskManager, null, testExecutor, testExecutor, TimeUnit.SECONDS.toNanos(1)); Task expectedTask = createSampleTask("null_msg_task", TaskState.TASK_STATE_WORKING, "ctx1"); when(mockTaskManager.getTask()).thenReturn(expectedTask); @@ -226,7 +226,7 @@ void testMultipleGetCurrentResultCalls() { void testGetCurrentResultWithMessageTakesPrecedence() { // Test that when both message and task are available, message takes precedence Message message = createSampleMessage("priority message", "pri1", Message.Role.ROLE_USER); - ResultAggregator messageAggregator = new ResultAggregator(mockTaskManager, message, testExecutor, testExecutor); + ResultAggregator messageAggregator = new ResultAggregator(mockTaskManager, message, testExecutor, testExecutor, TimeUnit.SECONDS.toNanos(1)); // Even if we set up the task manager to return something, message should take precedence Task task = createSampleTask("should_not_be_returned", TaskState.TASK_STATE_WORKING, "ctx1"); @@ -297,7 +297,7 @@ void testBlockingReturnsTaskWhenStoreBecomesVisibleAfterEmptyCapture() throws Ex TaskManager taskManager = new TaskManager(taskId, "ctx1", taskStore, null); ResultAggregator blockingAggregator = - new ResultAggregator(taskManager, null, testExecutor, testExecutor); + new ResultAggregator(taskManager, null, testExecutor, testExecutor, TimeUnit.SECONDS.toNanos(1)); MainEventBus mainEventBus = new MainEventBus(); InMemoryQueueManager queueManager = @@ -325,7 +325,7 @@ void testBlockingStillRaisesMissingTaskErrorWhenStoreRemainsEmpty() { TaskStore taskStore = mock(TaskStore.class); TaskManager taskManager = new TaskManager(taskId, "ctx1", taskStore, null); ResultAggregator blockingAggregator = - new ResultAggregator(taskManager, null, testExecutor, testExecutor); + new ResultAggregator(taskManager, null, testExecutor, testExecutor, TimeUnit.MILLISECONDS.toNanos(100)); InternalError error = assertThrows(InternalError.class, () -> blockingAggregator.consumeAndBreakOnInterrupt(createClosedEventConsumer(taskId), true)); @@ -339,7 +339,7 @@ void testNonBlockingDoesNotPollTaskStoreWhenCaptureIsEmpty() { TaskStore taskStore = mock(TaskStore.class); TaskManager taskManager = new TaskManager(taskId, "ctx1", taskStore, null); ResultAggregator nonBlockingAggregator = - new ResultAggregator(taskManager, null, testExecutor, testExecutor); + new ResultAggregator(taskManager, null, testExecutor, testExecutor, TimeUnit.SECONDS.toNanos(1)); assertThrows(InternalError.class, () -> nonBlockingAggregator.consumeAndBreakOnInterrupt(createClosedEventConsumer(taskId), false)); From 05db8a7efdbb00128ef491fe8f3e14f03369962e Mon Sep 17 00:00:00 2001 From: King Star Date: Thu, 23 Jul 2026 20:45:04 +0800 Subject: [PATCH 3/4] fix(server): preserve ResultAggregator constructor compatibility Keep the existing four-argument constructor delegating to the original one-second reconciliation timeout, and cover the new configuration wiring and default value. This fixes #2049 --- .../MicroProfileConfigProviderTest.java | 8 ++++++++ .../sdk/server/tasks/ResultAggregator.java | 7 +++++++ .../DefaultRequestHandlerTest.java | 19 +++++++++++++++++++ .../server/tasks/ResultAggregatorTest.java | 2 +- 4 files changed, 35 insertions(+), 1 deletion(-) diff --git a/integrations/microprofile-config/src/test/java/org/a2aproject/sdk/integrations/microprofile/MicroProfileConfigProviderTest.java b/integrations/microprofile-config/src/test/java/org/a2aproject/sdk/integrations/microprofile/MicroProfileConfigProviderTest.java index 2d85911ac..3fc0e429f 100644 --- a/integrations/microprofile-config/src/test/java/org/a2aproject/sdk/integrations/microprofile/MicroProfileConfigProviderTest.java +++ b/integrations/microprofile-config/src/test/java/org/a2aproject/sdk/integrations/microprofile/MicroProfileConfigProviderTest.java @@ -24,6 +24,8 @@ public class MicroProfileConfigProviderTest { private static final String A2A_EXECUTOR_CORE_POOL_SIZE = "a2a.executor.core-pool-size"; private static final String A2A_EXECUTOR_MAX_POOL_SIZE = "a2a.executor.max-pool-size"; private static final String A2A_EXECUTOR_KEEP_ALIVE_SECONDS = "a2a.executor.keep-alive-seconds"; + private static final String A2A_BLOCKING_RECONCILIATION_TIMEOUT_SECONDS = + "a2a.blocking.reconciliation.timeout.seconds"; @Inject A2AConfigProvider configProvider; @@ -59,6 +61,12 @@ public void testGetValueAnotherDefault() { assertEquals("60", value, "Should fall back to default value"); } + @Test + public void testGetReconciliationTimeoutDefault() { + String value = configProvider.getValue(A2A_BLOCKING_RECONCILIATION_TIMEOUT_SECONDS); + assertEquals("1", value, "Should fall back to the reconciliation timeout default"); + } + @Test public void testGetOptionalValueFromMicroProfileConfig() { // Test optional value that exists in application.properties diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java b/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java index 32c0b4ad0..b9a3c18a6 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java @@ -29,6 +29,7 @@ public class ResultAggregator { private static final Logger LOGGER = LoggerFactory.getLogger(ResultAggregator.class); + private static final long DEFAULT_TASK_STORE_RECONCILIATION_TIMEOUT_NANOS = TimeUnit.SECONDS.toNanos(1); private static final long TASK_STORE_RECONCILIATION_POLL_MILLIS = 10; private final TaskManager taskManager; @@ -37,6 +38,12 @@ public class ResultAggregator { private final long reconciliationTimeoutNanos; private volatile @Nullable Message message; + public ResultAggregator(TaskManager taskManager, @Nullable Message message, Executor executor, + Executor eventConsumerExecutor) { + this(taskManager, message, executor, eventConsumerExecutor, + DEFAULT_TASK_STORE_RECONCILIATION_TIMEOUT_NANOS); + } + public ResultAggregator(TaskManager taskManager, @Nullable Message message, Executor executor, Executor eventConsumerExecutor, long reconciliationTimeoutNanos) { this.taskManager = taskManager; diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/requesthandlers/DefaultRequestHandlerTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/requesthandlers/DefaultRequestHandlerTest.java index be84073ba..85ba534e3 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/requesthandlers/DefaultRequestHandlerTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/requesthandlers/DefaultRequestHandlerTest.java @@ -6,6 +6,8 @@ import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; import java.util.List; import java.util.Map; @@ -20,6 +22,7 @@ import org.a2aproject.sdk.server.ServerCallContext; import org.a2aproject.sdk.server.agentexecution.AgentExecutor; import org.a2aproject.sdk.server.agentexecution.RequestContext; +import org.a2aproject.sdk.server.config.A2AConfigProvider; import org.a2aproject.sdk.server.events.EventQueue; import org.a2aproject.sdk.server.events.EventQueueItem; import org.a2aproject.sdk.server.events.EventQueueUtil; @@ -146,6 +149,22 @@ protected interface AgentExecutorMethod { void invoke(RequestContext context, AgentEmitter agentEmitter) throws A2AError; } + @Test + void testInitConfigReadsBlockingTimeouts() { + A2AConfigProvider configProvider = mock(A2AConfigProvider.class); + when(configProvider.getValue("a2a.blocking.agent.timeout.seconds")).thenReturn("30"); + when(configProvider.getValue("a2a.blocking.consumption.timeout.seconds")).thenReturn("5"); + when(configProvider.getValue("a2a.blocking.reconciliation.timeout.seconds")).thenReturn("7"); + + DefaultRequestHandler handler = new DefaultRequestHandler(); + handler.configProvider = configProvider; + handler.initConfig(); + + assertEquals(30, handler.agentCompletionTimeoutSeconds); + assertEquals(5, handler.consumptionCompletionTimeoutSeconds); + assertEquals(7, handler.reconciliationTimeoutSeconds); + } + /** * Test 1: Non-streaming AUTH_REQUIRED returns immediately while agent continues. * Verifies: diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java index d3adc986e..91d4b50c6 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java @@ -62,7 +62,7 @@ public class ResultAggregatorTest { @BeforeEach void setUp() { MockitoAnnotations.openMocks(this); - aggregator = new ResultAggregator(mockTaskManager, null, testExecutor, testExecutor, TimeUnit.SECONDS.toNanos(1)); + aggregator = new ResultAggregator(mockTaskManager, null, testExecutor, testExecutor); } // Helper methods for creating sample data From 0239faacb0cae95303bf2ca8ed30a0a996fbd1c0 Mon Sep 17 00:00:00 2001 From: King Star Date: Thu, 23 Jul 2026 23:06:30 +0800 Subject: [PATCH 4/4] fix(server): remove unused ResultAggregator overload --- .../org/a2aproject/sdk/server/tasks/ResultAggregator.java | 7 ------- .../a2aproject/sdk/server/tasks/ResultAggregatorTest.java | 3 ++- 2 files changed, 2 insertions(+), 8 deletions(-) diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java b/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java index b9a3c18a6..32c0b4ad0 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/tasks/ResultAggregator.java @@ -29,7 +29,6 @@ public class ResultAggregator { private static final Logger LOGGER = LoggerFactory.getLogger(ResultAggregator.class); - private static final long DEFAULT_TASK_STORE_RECONCILIATION_TIMEOUT_NANOS = TimeUnit.SECONDS.toNanos(1); private static final long TASK_STORE_RECONCILIATION_POLL_MILLIS = 10; private final TaskManager taskManager; @@ -38,12 +37,6 @@ public class ResultAggregator { private final long reconciliationTimeoutNanos; private volatile @Nullable Message message; - public ResultAggregator(TaskManager taskManager, @Nullable Message message, Executor executor, - Executor eventConsumerExecutor) { - this(taskManager, message, executor, eventConsumerExecutor, - DEFAULT_TASK_STORE_RECONCILIATION_TIMEOUT_NANOS); - } - public ResultAggregator(TaskManager taskManager, @Nullable Message message, Executor executor, Executor eventConsumerExecutor, long reconciliationTimeoutNanos) { this.taskManager = taskManager; diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java index 91d4b50c6..04c6d582a 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/tasks/ResultAggregatorTest.java @@ -62,7 +62,8 @@ public class ResultAggregatorTest { @BeforeEach void setUp() { MockitoAnnotations.openMocks(this); - aggregator = new ResultAggregator(mockTaskManager, null, testExecutor, testExecutor); + aggregator = new ResultAggregator(mockTaskManager, null, testExecutor, testExecutor, + TimeUnit.SECONDS.toNanos(1)); } // Helper methods for creating sample data