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/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/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 eddc57bbe..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
@@ -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,17 +29,21 @@
public class ResultAggregator {
private static final Logger LOGGER = LoggerFactory.getLogger(ResultAggregator.class);
+ 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() {
@@ -235,7 +240,8 @@ 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 poll TaskStore with bounded timeout (blocking) or single read (non-blocking).
EventKind eventKind = message.get();
if (eventKind == null) {
eventKind = capturedTask.get();
@@ -244,7 +250,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 +265,33 @@ else if (blocking) {
consumptionCompletionFuture);
}
+ private @Nullable Task reconcileTaskStore(boolean blocking) throws A2AError {
+ Task task = taskManager.getTask();
+ if (task != null || !blocking) {
+ return task;
+ }
+
+ 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);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new InternalError("Interrupted while reconciling TaskStore for " + taskManager.getTaskId());
+ }
+
+ task = taskManager.getTask();
+ if (task != null) {
+ LOGGER.debug("TaskStore reconciliation: task {} found after polling", task.id());
+ return task;
+ }
+ }
+ return null;
+ }
+
private String taskIdForLogging() {
Task task = taskManager.getTask();
return task != null ? task.id() : "unknown";
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/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 7b206e01a..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
@@ -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;
@@ -57,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
@@ -115,7 +121,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());
@@ -126,7 +132,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();
@@ -161,7 +167,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);
@@ -221,7 +227,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");
@@ -281,6 +287,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, TimeUnit.SECONDS.toNanos(1));
+
+ 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, TimeUnit.MILLISECONDS.toNanos(100));
+
+ 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, TimeUnit.SECONDS.toNanos(1));
+
+ 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