From d58f5a8eb0755bd3dc125e1949a55930d420e346 Mon Sep 17 00:00:00 2001 From: 014-code <2402143478@qq.com> Date: Sat, 25 Jul 2026 22:37:02 +0800 Subject: [PATCH 1/3] fix(server): start event bus processor for manual wiring Ensure MainEventBusProcessor.ensureStarted() starts manually constructed processors, and call it from DefaultRequestHandler.create(...). This prevents manually wired integrations from leaving events stuck on the main event bus. Add a regression test for the manual wiring path that returns a Message response. This fixes #995 --- .../server/events/MainEventBusProcessor.java | 18 +++++++--- .../DefaultRequestHandler.java | 1 + .../DefaultRequestHandlerTest.java | 35 +++++++++++++++++++ 3 files changed, 49 insertions(+), 5 deletions(-) diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/events/MainEventBusProcessor.java b/server-common/src/main/java/org/a2aproject/sdk/server/events/MainEventBusProcessor.java index eab991ff7..72f5f9204 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/events/MainEventBusProcessor.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/events/MainEventBusProcessor.java @@ -137,7 +137,12 @@ public void setPushNotificationExecutor(java.util.concurrent.Executor executor) @SuppressWarnings("NullAway.Init") @PostConstruct - void start() { + synchronized void start() { + if (processorThread != null) { + LOGGER.debug("MainEventBusProcessor already started"); + return; + } + running = true; processorThread = new Thread(this, "MainEventBusProcessor"); processorThread.setDaemon(true); // Allow JVM to exit even if this thread is running processorThread.start(); @@ -145,15 +150,18 @@ void start() { } /** - * No-op method to force CDI proxy resolution and ensure @PostConstruct has been called. - * Called by MainEventBusProcessorInitializer during application startup. + * Ensures the background processor thread has been started. + *
+ * In CDI runtimes, this forces proxy resolution and {@link PostConstruct}. For manual + * wiring, this starts the processor directly. + *
*/ public void ensureStarted() { - // Method intentionally empty - just forces proxy resolution + start(); } @PreDestroy - void stop() { + synchronized void stop() { LOGGER.info("MainEventBusProcessor stopping..."); running = false; if (processorThread != null) { 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 addac8e69..67671650b 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 @@ -309,6 +309,7 @@ public static DefaultRequestHandler create(AgentExecutor agentExecutor, TaskStor handler.agentCompletionTimeoutSeconds = 5; handler.consumptionCompletionTimeoutSeconds = 2; handler.reconciliationTimeoutSeconds = 1; + handler.mainEventBusProcessor.ensureStarted(); return handler; } 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 85ba534e3..547b54706 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 @@ -165,6 +165,41 @@ void testInitConfigReadsBlockingTimeouts() { assertEquals(7, handler.reconciliationTimeoutSeconds); } + @Test + void testCreateStartsManuallyConstructedMainEventBusProcessor() throws Exception { + InMemoryTaskStore manualTaskStore = new InMemoryTaskStore(); + PushNotificationConfigStore manualPushConfigStore = new InMemoryPushNotificationConfigStore(); + MainEventBus manualMainEventBus = new MainEventBus(); + InMemoryQueueManager manualQueueManager = new InMemoryQueueManager(manualTaskStore, manualMainEventBus); + MainEventBusProcessor manualProcessor = new MainEventBusProcessor( + manualMainEventBus, manualTaskStore, NOOP_PUSHNOTIFICATION_SENDER, manualQueueManager); + + try { + AgentExecutor manualAgentExecutor = new AgentExecutor() { + @Override + public void execute(RequestContext context, AgentEmitter agentEmitter) { + agentEmitter.sendMessage("manual lifecycle response"); + } + + @Override + public void cancel(RequestContext context, AgentEmitter agentEmitter) { + throw new AssertionError("Cancel should not be invoked"); + } + }; + RequestHandler manualRequestHandler = DefaultRequestHandler.create( + manualAgentExecutor, manualTaskStore, manualQueueManager, manualPushConfigStore, + manualProcessor, internalExecutor, internalExecutor); + + EventKind eventKind = manualRequestHandler.onMessageSend( + MessageSendParams.builder().message(MESSAGE).build(), NULL_CONTEXT); + + assertInstanceOf(Message.class, eventKind); + assertEquals(Message.Role.ROLE_AGENT, ((Message) eventKind).role()); + } finally { + EventQueueUtil.stop(manualProcessor); + } + } + /** * Test 1: Non-streaming AUTH_REQUIRED returns immediately while agent continues. * Verifies: From ea0b91d140b02dbb9f1271b14eedfbc0ade36ebb Mon Sep 17 00:00:00 2001 From: 014-code <2402143478@qq.com> Date: Mon, 27 Jul 2026 20:43:17 +0800 Subject: [PATCH 2/3] fix(server): cover direct request handler wiring Start the main event bus processor from the public request handler constructor and verify the direct manual wiring path.\n\nThis fixes #995 --- .../sdk/server/events/MainEventBusProcessor.java | 2 +- .../sdk/server/requesthandlers/DefaultRequestHandler.java | 2 +- .../server/requesthandlers/DefaultRequestHandlerTest.java | 7 +++++-- 3 files changed, 7 insertions(+), 4 deletions(-) diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/events/MainEventBusProcessor.java b/server-common/src/main/java/org/a2aproject/sdk/server/events/MainEventBusProcessor.java index 72f5f9204..c8fe8f027 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/events/MainEventBusProcessor.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/events/MainEventBusProcessor.java @@ -138,7 +138,7 @@ public void setPushNotificationExecutor(java.util.concurrent.Executor executor) @SuppressWarnings("NullAway.Init") @PostConstruct synchronized void start() { - if (processorThread != null) { + if (processorThread != null && processorThread.isAlive()) { LOGGER.debug("MainEventBusProcessor already started"); return; } 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 67671650b..6d8de00b7 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 @@ -282,6 +282,7 @@ public DefaultRequestHandler(AgentExecutor agentExecutor, TaskStore taskStore, // I am unsure about the correct scope. // Also reworked to make a Supplier since otherwise the builder gets polluted with wrong tasks this.requestContextBuilder = () -> new SimpleRequestContextBuilder(taskStore, false); + this.mainEventBusProcessor.ensureStarted(); } @SuppressWarnings("NullAway.Init") @@ -309,7 +310,6 @@ public static DefaultRequestHandler create(AgentExecutor agentExecutor, TaskStor handler.agentCompletionTimeoutSeconds = 5; handler.consumptionCompletionTimeoutSeconds = 2; handler.reconciliationTimeoutSeconds = 1; - handler.mainEventBusProcessor.ensureStarted(); return handler; } 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 547b54706..936369e98 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 @@ -166,7 +166,7 @@ void testInitConfigReadsBlockingTimeouts() { } @Test - void testCreateStartsManuallyConstructedMainEventBusProcessor() throws Exception { + void testConstructorStartsManuallyConstructedMainEventBusProcessor() throws Exception { InMemoryTaskStore manualTaskStore = new InMemoryTaskStore(); PushNotificationConfigStore manualPushConfigStore = new InMemoryPushNotificationConfigStore(); MainEventBus manualMainEventBus = new MainEventBus(); @@ -186,9 +186,12 @@ public void cancel(RequestContext context, AgentEmitter agentEmitter) { throw new AssertionError("Cancel should not be invoked"); } }; - RequestHandler manualRequestHandler = DefaultRequestHandler.create( + DefaultRequestHandler manualRequestHandler = new DefaultRequestHandler( manualAgentExecutor, manualTaskStore, manualQueueManager, manualPushConfigStore, manualProcessor, internalExecutor, internalExecutor); + manualRequestHandler.agentCompletionTimeoutSeconds = 5; + manualRequestHandler.consumptionCompletionTimeoutSeconds = 2; + manualRequestHandler.reconciliationTimeoutSeconds = 1; EventKind eventKind = manualRequestHandler.onMessageSend( MessageSendParams.builder().message(MESSAGE).build(), NULL_CONTEXT); From 6a598233084547adcdf932f67521f8e8aa40fe2a Mon Sep 17 00:00:00 2001 From: 014-code <2402143478@qq.com> Date: Mon, 27 Jul 2026 21:54:40 +0800 Subject: [PATCH 3/3] fix(server): preserve task ids for emitted messages Ensure AgentEmitter fills missing task and context IDs when sending an existing Message, so streaming clients can reliably observe the server-generated task ID for message-only responses. This fixes #995 --- .../a2aproject/sdk/server/tasks/AgentEmitter.java | 12 ++++++++++-- .../sdk/server/tasks/AgentEmitterTest.java | 6 ++++++ 2 files changed, 16 insertions(+), 2 deletions(-) diff --git a/server-common/src/main/java/org/a2aproject/sdk/server/tasks/AgentEmitter.java b/server-common/src/main/java/org/a2aproject/sdk/server/tasks/AgentEmitter.java index 19d458fd6..e31b6a471 100644 --- a/server-common/src/main/java/org/a2aproject/sdk/server/tasks/AgentEmitter.java +++ b/server-common/src/main/java/org/a2aproject/sdk/server/tasks/AgentEmitter.java @@ -491,7 +491,8 @@ public void sendMessage(List* Use this when you need to forward or echo an existing message without creating a new one. - * The message is enqueued as-is, preserving its messageId, metadata, and all other fields. + * The message is enqueued with this emitter's task and context IDs when they are missing, + * preserving its messageId, metadata, and all other fields. *
** Note: This is typically used for forwarding user messages or preserving specific @@ -518,7 +519,14 @@ public void sendMessage(Message message) { LOGGER.error("Message contextId mismatch: expected={}, actual={}", contextId, message.contextId()); throw new IllegalArgumentException("Message contextId does not match the emitter's contextId"); } - eventQueue.enqueueEvent(message); + Message messageToSend = message; + if (message.taskId() == null || message.contextId() == null) { + messageToSend = Message.builder(message) + .taskId(taskId) + .contextId(contextId) + .build(); + } + eventQueue.enqueueEvent(messageToSend); } /** diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/tasks/AgentEmitterTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/tasks/AgentEmitterTest.java index 233d49bad..e95c63c53 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/tasks/AgentEmitterTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/tasks/AgentEmitterTest.java @@ -461,6 +461,12 @@ public void sendMessageWithNullIdsSucceeds() throws Exception { EventQueueItem item = eventQueue.dequeueEventItem(WAIT_MILLI_SECONDS); assertNotNull(item); assertInstanceOf(Message.class, item.getEvent()); + Message emittedMessage = (Message) item.getEvent(); + assertEquals(TEST_TASK_ID, emittedMessage.taskId()); + assertEquals(TEST_TASK_CONTEXT_ID, emittedMessage.contextId()); + assertEquals(message.messageId(), emittedMessage.messageId()); + assertEquals(message.role(), emittedMessage.role()); + assertEquals(message.parts(), emittedMessage.parts()); } @Test