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..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 @@ -137,7 +137,12 @@ public void setPushNotificationExecutor(java.util.concurrent.Executor executor) @SuppressWarnings("NullAway.Init") @PostConstruct - void start() { + synchronized void start() { + if (processorThread != null && processorThread.isAlive()) { + 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..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") 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> parts, @Nullable Map metad * Sends an existing Message object directly to the client. *

* 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/requesthandlers/DefaultRequestHandlerTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/requesthandlers/DefaultRequestHandlerTest.java index 85ba534e3..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 @@ -165,6 +165,44 @@ void testInitConfigReadsBlockingTimeouts() { assertEquals(7, handler.reconciliationTimeoutSeconds); } + @Test + void testConstructorStartsManuallyConstructedMainEventBusProcessor() 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"); + } + }; + 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); + + 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: 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