From 864939a49f447b9577874ef683cd3b16ea4f5bf6 Mon Sep 17 00:00:00 2001 From: xiaoming-2026 Date: Sat, 8 Aug 2026 02:28:05 +0800 Subject: [PATCH 1/9] feat(a2a): publish remote agent catalog snapshots --- .../autoconfigure/A2AAutoConfiguration.java | 5 +- .../client/A2ARemoteAgentCardRegistry.java | 78 ++++++++++++- .../RemoteAgentCatalogChangedEvent.java | 14 +++ .../client/RemoteAgentCatalogSnapshot.java | 23 ++++ .../A2AAutoConfigurationTest.java | 29 +++++ .../A2ARemoteAgentCardRegistryTest.java | 106 ++++++++++++++++++ 6 files changed, 248 insertions(+), 7 deletions(-) create mode 100644 service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogChangedEvent.java create mode 100644 service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogSnapshot.java create mode 100644 service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistryTest.java diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java index bf8882d4..ef07b9e0 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java @@ -50,6 +50,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.annotation.Bean; import org.springframework.core.env.Environment; @@ -237,8 +238,8 @@ public A2AAgentExecutor a2aAgentExecutor(ServeOrchestrator orchestrator, A2AProt */ @Bean @ConditionalOnMissingBean - public A2ARemoteAgentCardRegistry a2aRemoteAgentCardRegistry() { - return new A2ARemoteAgentCardRegistry(); + public A2ARemoteAgentCardRegistry a2aRemoteAgentCardRegistry(ApplicationEventPublisher eventPublisher) { + return new A2ARemoteAgentCardRegistry(eventPublisher); } /** diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistry.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistry.java index 0acb76ad..224c43fe 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistry.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistry.java @@ -5,26 +5,54 @@ package com.openjiuwen.service.app.controller.a2a.client; import org.a2aproject.sdk.spec.AgentCard; -import org.springframework.stereotype.Component; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.context.ApplicationEventPublisher; import java.util.List; import java.util.Map; import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.locks.ReentrantLock; /** * Thread-safe in-memory registry of discovered remote A2A AgentCards. * * @since 0.1.0 */ -@Component public class A2ARemoteAgentCardRegistry { + private static final Logger log = LoggerFactory.getLogger(A2ARemoteAgentCardRegistry.class); + /** * Default timeout in seconds for remote agent calls. */ static final int DEFAULT_TIMEOUT_SECONDS = 300; + private final ApplicationEventPublisher eventPublisher; private final Map entries = new ConcurrentHashMap<>(); + private final ReentrantLock updateLock = new ReentrantLock(); + + private long version; + + /** + * Creates a registry without event publication. + * + *

This constructor preserves direct, non-Spring usage. Runtime auto-configuration + * supplies an {@link ApplicationEventPublisher}.

+ */ + public A2ARemoteAgentCardRegistry() { + this(event -> { + }); + } + + /** + * Creates a registry that publishes complete catalog snapshots after updates. + * + * @param eventPublisher the Spring application event publisher + */ + public A2ARemoteAgentCardRegistry(ApplicationEventPublisher eventPublisher) { + this.eventPublisher = eventPublisher; + } /** * Registers a remote agent card using the default timeout. @@ -42,7 +70,7 @@ public void register(String name, AgentCard card) { * @return an unmodifiable copy of all entries */ public List getAll() { - return List.copyOf(entries.values()); + return snapshot().entries(); } /** @@ -73,10 +101,25 @@ public String resolveUrl(String name) { return ifaces.get(0).url(); } + /** + * Returns the current complete remote-agent catalog. + * + * @return an immutable, name-sorted catalog snapshot + */ + public RemoteAgentCatalogSnapshot snapshot() { + updateLock.lock(); + try { + return createSnapshot(); + } finally { + updateLock.unlock(); + } + } + /** * A registered remote agent entry holding the card and timeout configuration. */ - public record RemoteAgentEntry(String name, AgentCard card, int timeoutSeconds, boolean isStreaming) {} + public record RemoteAgentEntry(String name, AgentCard card, int timeoutSeconds, boolean isStreaming) { + } /** * Registers a remote agent card with a specific timeout. @@ -98,6 +141,31 @@ public void register(String name, AgentCard card, int timeoutSeconds) { * @param isStreaming whether Runtime should prefer a streaming remote invocation */ public void register(String name, AgentCard card, int timeoutSeconds, boolean isStreaming) { - entries.put(name, new RemoteAgentEntry(name, card, timeoutSeconds, isStreaming)); + RemoteAgentCatalogSnapshot updatedSnapshot; + updateLock.lock(); + try { + entries.put(name, new RemoteAgentEntry(name, card, timeoutSeconds, isStreaming)); + version++; + updatedSnapshot = createSnapshot(); + } finally { + updateLock.unlock(); + } + publishCatalogChanged(updatedSnapshot); + } + + private RemoteAgentCatalogSnapshot createSnapshot() { + List sortedEntries = entries.values().stream() + .sorted((left, right) -> left.name().compareTo(right.name())).toList(); + return new RemoteAgentCatalogSnapshot(version, sortedEntries); + } + + private void publishCatalogChanged(RemoteAgentCatalogSnapshot updatedSnapshot) { + try { + eventPublisher.publishEvent(new RemoteAgentCatalogChangedEvent(updatedSnapshot)); + } catch (RuntimeException exception) { + log.error("Failed to publish remote Agent Card catalog event, version={}", updatedSnapshot.version(), + exception); + throw exception; + } } } diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogChangedEvent.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogChangedEvent.java new file mode 100644 index 00000000..04aa6abf --- /dev/null +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogChangedEvent.java @@ -0,0 +1,14 @@ +/* + * Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved. + */ + +package com.openjiuwen.service.app.controller.a2a.client; + +/** + * Event published after the remote A2A agent catalog changes. + * + * @param snapshot complete catalog snapshot produced by the registry update + * @since 0.1.1 + */ +public record RemoteAgentCatalogChangedEvent(RemoteAgentCatalogSnapshot snapshot) { +} diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogSnapshot.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogSnapshot.java new file mode 100644 index 00000000..11b22bf4 --- /dev/null +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogSnapshot.java @@ -0,0 +1,23 @@ +/* + * Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved. + */ + +package com.openjiuwen.service.app.controller.a2a.client; + +import java.util.List; + +/** + * Immutable versioned snapshot of all discovered remote A2A agents. + * + * @param version monotonically increasing registry version + * @param entries complete remote-agent entries sorted by name + * @since 0.1.1 + */ +public record RemoteAgentCatalogSnapshot(long version, List entries) { + /** + * Creates an immutable snapshot. + */ + public RemoteAgentCatalogSnapshot { + entries = List.copyOf(entries); + } +} diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfigurationTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfigurationTest.java index 7b275a70..8f18956f 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfigurationTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfigurationTest.java @@ -8,13 +8,20 @@ import static org.mockito.Mockito.mock; import com.openjiuwen.service.app.config.SpringEnvironmentConfigProvider; +import com.openjiuwen.service.app.controller.a2a.client.A2ARemoteAgentCardRegistry; +import com.openjiuwen.service.app.controller.a2a.client.RemoteAgentCatalogChangedEvent; import com.openjiuwen.service.spec.spi.ServeOrchestrator; import org.a2aproject.sdk.server.config.A2AConfigProvider; import org.a2aproject.sdk.server.requesthandlers.RequestHandler; +import org.a2aproject.sdk.spec.AgentCard; import org.junit.jupiter.api.Test; import org.springframework.boot.autoconfigure.AutoConfigurations; import org.springframework.boot.test.context.runner.ApplicationContextRunner; +import org.springframework.context.PayloadApplicationEvent; + +import java.util.ArrayList; +import java.util.List; /** * Auto-configuration tests for A2A SDK configuration. @@ -68,4 +75,26 @@ void a2aConfigProviderAllowsCustomProviderOverride() { assertThat(context.getBean(A2AConfigProvider.class)).isSameAs(customProvider); }); } + + @Test + void remoteAgentCardRegistryPublishesCatalogChanges() { + contextRunner.run(context -> { + List events = new ArrayList<>(); + context.getSourceApplicationContext().addApplicationListener(event -> { + if (event instanceof PayloadApplicationEvent payloadEvent + && payloadEvent.getPayload() instanceof RemoteAgentCatalogChangedEvent catalogChangedEvent) { + events.add(catalogChangedEvent); + } + }); + + A2ARemoteAgentCardRegistry registry = context.getBean(A2ARemoteAgentCardRegistry.class); + registry.register("balance", mock(AgentCard.class)); + + assertThat(events).singleElement().satisfies(event -> { + assertThat(event.snapshot().version()).isEqualTo(1L); + assertThat(event.snapshot().entries()).extracting(A2ARemoteAgentCardRegistry.RemoteAgentEntry::name) + .containsExactly("balance"); + }); + }); + } } diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistryTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistryTest.java new file mode 100644 index 00000000..fe7b6f15 --- /dev/null +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistryTest.java @@ -0,0 +1,106 @@ +/* + * Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved. + */ + +package com.openjiuwen.service.app.controller.a2a.client; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.mock; + +import org.a2aproject.sdk.spec.AgentCard; +import org.junit.jupiter.api.Test; +import org.springframework.context.ApplicationEventPublisher; + +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.stream.IntStream; + +class A2ARemoteAgentCardRegistryTest { + @Test + void initialSnapshotIsEmptyAndImmutable() { + A2ARemoteAgentCardRegistry registry = new A2ARemoteAgentCardRegistry(); + + RemoteAgentCatalogSnapshot snapshot = registry.snapshot(); + + assertThat(snapshot.version()).isZero(); + assertThat(snapshot.entries()).isEmpty(); + assertThatThrownBy(() -> snapshot.entries().add(entry("other"))) + .isInstanceOf(UnsupportedOperationException.class); + } + + @Test + void registrationPublishesCompleteSortedSnapshots() { + List events = new CopyOnWriteArrayList<>(); + A2ARemoteAgentCardRegistry registry = registryWithEvents(events); + + registry.register("weather", mock(AgentCard.class), 30, true); + registry.register("balance", mock(AgentCard.class), 60, false); + + assertThat(events).hasSize(2); + assertThat(events.get(0).snapshot().version()).isEqualTo(1L); + assertThat(events.get(0).snapshot().entries()).extracting(A2ARemoteAgentCardRegistry.RemoteAgentEntry::name) + .containsExactly("weather"); + assertThat(events.get(1).snapshot().version()).isEqualTo(2L); + assertThat(events.get(1).snapshot().entries()).extracting(A2ARemoteAgentCardRegistry.RemoteAgentEntry::name) + .containsExactly("balance", "weather"); + assertThat(registry.getAll()).containsExactlyElementsOf(events.get(1).snapshot().entries()); + } + + @Test + void replacingSameNameCreatesNewVersion() { + List events = new CopyOnWriteArrayList<>(); + A2ARemoteAgentCardRegistry registry = registryWithEvents(events); + AgentCard firstCard = mock(AgentCard.class); + AgentCard secondCard = mock(AgentCard.class); + + registry.register("transfer", firstCard, 30, false); + registry.register("transfer", secondCard, 90, true); + + assertThat(registry.snapshot().version()).isEqualTo(2L); + assertThat(registry.snapshot().entries()).singleElement().satisfies(entry -> { + assertThat(entry.card()).isSameAs(secondCard); + assertThat(entry.timeoutSeconds()).isEqualTo(90); + assertThat(entry.isStreaming()).isTrue(); + }); + assertThat(events).extracting(event -> event.snapshot().version()).containsExactly(1L, 2L); + } + + @Test + void concurrentRegistrationProducesUniqueCompleteVersions() { + List events = new CopyOnWriteArrayList<>(); + A2ARemoteAgentCardRegistry registry = registryWithEvents(events); + + IntStream.range(0, 32).parallel() + .forEach(index -> registry.register("agent-" + index, mock(AgentCard.class), 30, false)); + + assertThat(registry.snapshot().version()).isEqualTo(32L); + assertThat(registry.snapshot().entries()).hasSize(32); + assertThat(events).hasSize(32); + assertThat(events).extracting(event -> event.snapshot().version()).doesNotHaveDuplicates() + .containsExactlyInAnyOrderElementsOf(IntStream.rangeClosed(1, 32).mapToObj(Long::valueOf).toList()); + assertThat(events) + .allSatisfy(event -> assertThat(event.snapshot().entries()).hasSize((int) event.snapshot().version())); + } + + @Test + void publicationFailureKeepsCompletedRegistryUpdateVisible() { + A2ARemoteAgentCardRegistry registry = new A2ARemoteAgentCardRegistry(event -> { + throw new IllegalStateException("listener failed"); + }); + + assertThatThrownBy(() -> registry.register("balance", mock(AgentCard.class))) + .isInstanceOf(IllegalStateException.class).hasMessage("listener failed"); + assertThat(registry.snapshot().version()).isEqualTo(1L); + assertThat(registry.get("balance")).isPresent(); + } + + private static A2ARemoteAgentCardRegistry registryWithEvents(List events) { + ApplicationEventPublisher publisher = event -> events.add((RemoteAgentCatalogChangedEvent) event); + return new A2ARemoteAgentCardRegistry(publisher); + } + + private static A2ARemoteAgentCardRegistry.RemoteAgentEntry entry(String name) { + return new A2ARemoteAgentCardRegistry.RemoteAgentEntry(name, mock(AgentCard.class), 30, false); + } +} From 35cee82085fa317bed3c8cb9e9c24062b75c4a70 Mon Sep 17 00:00:00 2001 From: xiaoming-2026 Date: Sat, 8 Aug 2026 03:23:57 +0800 Subject: [PATCH 2/9] fix(agentcore): preserve DeepAgent interruptions --- .../adapters/agentcore/agentfw/JiuwenCoreAgentHandler.java | 7 ++++--- .../agentcore/agentfw/JiuwenCoreAgentHandlerTest.java | 7 +++++++ 2 files changed, 11 insertions(+), 3 deletions(-) diff --git a/service/agent-service-adapters/agent-service-adapters-agentcore/src/main/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandler.java b/service/agent-service-adapters/agent-service-adapters-agentcore/src/main/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandler.java index 70086d01..46a88caa 100644 --- a/service/agent-service-adapters/agent-service-adapters-agentcore/src/main/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandler.java +++ b/service/agent-service-adapters/agent-service-adapters-agentcore/src/main/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandler.java @@ -370,10 +370,11 @@ private static QueryResponse buildQueryResponse(Object lastPayload, StringBuilde return new QueryResponse(result, conversationId); } - private static boolean supportsInvoke(Object agent) { - if (agent == null || agent instanceof String) { + static boolean supportsInvoke(Object agent) { + if (agent == null || agent instanceof String || agent instanceof DeepAgent) { // Resolved at runtime from agent-id; use streaming unless the instance exposes - // invoke. + // invoke. DeepAgent task-loop interruptions are emitted on its stream and may + // not be represented by the aggregate invoke result after an earlier tool call. return false; } for (Method method : agent.getClass().getMethods()) { diff --git a/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java b/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java index 6264656a..5a6a45b0 100644 --- a/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java +++ b/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java @@ -26,6 +26,7 @@ import com.openjiuwen.core.singleagent.interrupt.ToolCallInterruptRequest; import com.openjiuwen.core.singleagent.schema.AgentCard; import com.openjiuwen.core.workflow.WorkflowOutput; +import com.openjiuwen.harness.deep_agent.DeepAgent; import com.openjiuwen.service.adapters.agentcore.external.ExternalSvcAdapterRegistrar; import com.openjiuwen.service.spec.dto.QueryChunk; import com.openjiuwen.service.spec.dto.QueryResponse; @@ -229,6 +230,12 @@ void syncQueryUsesInvokePathWhenAgentSupportsIt() { assertThat((Map) second.getResult()).containsEntry("content", "turn2:b|prev=a"); } + @Test + void deepAgentUsesStreamingPathToPreserveTaskLoopInterruptions() { + assertThat(JiuwenCoreAgentHandler.supportsInvoke(mock(DeepAgent.class))).isFalse(); + assertThat(JiuwenCoreAgentHandler.supportsInvoke(new InvokeEchoAgent())).isTrue(); + } + @Test @SuppressWarnings("unchecked") void nonStreamingQueryPreservesAllRemoteInterruptsInOriginalOrder() { From 6f07ac69cedfb35f54489bb974227085265c2a67 Mon Sep 17 00:00:00 2001 From: xiaoming-2026 Date: Sat, 8 Aug 2026 04:02:14 +0800 Subject: [PATCH 3/9] fix(runtime): address intent routing quality checks --- .../agentfw/JiuwenCoreAgentHandlerTest.java | 2 +- .../app/autoconfigure/A2AAutoConfiguration.java | 1 + .../a2a/client/A2ARemoteAgentCardRegistry.java | 12 +----------- .../a2a/client/A2ARemoteAgentCardRegistryTest.java | 9 ++++++++- 4 files changed, 11 insertions(+), 13 deletions(-) diff --git a/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java b/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java index 5a6a45b0..45440680 100644 --- a/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java +++ b/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java @@ -45,7 +45,7 @@ import java.util.concurrent.atomic.AtomicReference; /** - * JiuwenCoreAgentHandlerTest + * Tests AgentCore request adaptation, aggregation, and interruption handling. * * @since 2026-07-03 */ diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java index ef07b9e0..51519990 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java @@ -234,6 +234,7 @@ public A2AAgentExecutor a2aAgentExecutor(ServeOrchestrator orchestrator, A2AProt /** * Creates the remote agent card registry bean. * + * @param eventPublisher the Spring application event publisher * @return the remote agent card registry */ @Bean diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistry.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistry.java index 224c43fe..60e0efe1 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistry.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistry.java @@ -5,8 +5,6 @@ package com.openjiuwen.service.app.controller.a2a.client; import org.a2aproject.sdk.spec.AgentCard; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import org.springframework.context.ApplicationEventPublisher; import java.util.List; @@ -21,8 +19,6 @@ * @since 0.1.0 */ public class A2ARemoteAgentCardRegistry { - private static final Logger log = LoggerFactory.getLogger(A2ARemoteAgentCardRegistry.class); - /** * Default timeout in seconds for remote agent calls. */ @@ -160,12 +156,6 @@ private RemoteAgentCatalogSnapshot createSnapshot() { } private void publishCatalogChanged(RemoteAgentCatalogSnapshot updatedSnapshot) { - try { - eventPublisher.publishEvent(new RemoteAgentCatalogChangedEvent(updatedSnapshot)); - } catch (RuntimeException exception) { - log.error("Failed to publish remote Agent Card catalog event, version={}", updatedSnapshot.version(), - exception); - throw exception; - } + eventPublisher.publishEvent(new RemoteAgentCatalogChangedEvent(updatedSnapshot)); } } diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistryTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistryTest.java index fe7b6f15..62c4e617 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistryTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistryTest.java @@ -16,6 +16,7 @@ import java.util.concurrent.CopyOnWriteArrayList; import java.util.stream.IntStream; +/** Tests remote Agent Card catalog snapshots and update publication. */ class A2ARemoteAgentCardRegistryTest { @Test void initialSnapshotIsEmptyAndImmutable() { @@ -96,7 +97,13 @@ void publicationFailureKeepsCompletedRegistryUpdateVisible() { } private static A2ARemoteAgentCardRegistry registryWithEvents(List events) { - ApplicationEventPublisher publisher = event -> events.add((RemoteAgentCatalogChangedEvent) event); + ApplicationEventPublisher publisher = event -> { + if (event instanceof RemoteAgentCatalogChangedEvent catalogEvent) { + events.add(catalogEvent); + return; + } + throw new IllegalArgumentException("Unexpected event type: " + event.getClass().getName()); + }; return new A2ARemoteAgentCardRegistry(publisher); } From b572c91994087b51b0c787508fce46455c72955e Mon Sep 17 00:00:00 2001 From: xiaoming-2026 Date: Sat, 8 Aug 2026 08:03:42 +0800 Subject: [PATCH 4/9] refactor(a2a): expose remote agent catalog API --- .../app/a2a/catalog.README.md" | 16 +++ .../app/controller/a2a/client.README.md" | 1 - .../catalog}/A2ARemoteAgentCardRegistry.java | 16 ++- .../RemoteAgentCatalogChangedEvent.java | 2 +- .../catalog}/RemoteAgentCatalogSnapshot.java | 4 +- .../app/a2a/catalog/RemoteAgentEntry.java | 19 +++ .../autoconfigure/A2AAutoConfiguration.java | 2 +- .../a2a/client/A2AAgentCardDiscovery.java | 1 + .../a2a/client/A2ARemoteAgentClient.java | 9 +- .../a2a/client/RemoteAgentCardResolver.java | 2 + .../A2ARemoteAgentCardRegistryTest.java | 10 +- .../A2AAutoConfigurationTest.java | 8 +- .../a2a/client/A2AAgentCardDiscoveryTest.java | 1 + .../A2ARemoteAgentClientClassLoaderTest.java | 68 +++++----- .../A2ARemoteAgentClientResultTest.java | 1 + ...moteAgentClientStreamingLifecycleTest.java | 1 + .../DualRuntimeCallbackIntegrationTest.java | 122 ++++++------------ .../it/DualRuntimeFailureIntegrationTest.java | 2 +- 18 files changed, 146 insertions(+), 139 deletions(-) create mode 100644 "documents/zh/2.\345\274\200\345\217\221\346\214\207\345\215\227/API\346\226\207\346\241\243/com.openjiuwen.service/app/a2a/catalog.README.md" rename service/agent-service-app/src/main/java/com/openjiuwen/service/app/{controller/a2a/client => a2a/catalog}/A2ARemoteAgentCardRegistry.java (88%) rename service/agent-service-app/src/main/java/com/openjiuwen/service/app/{controller/a2a/client => a2a/catalog}/RemoteAgentCatalogChangedEvent.java (85%) rename service/agent-service-app/src/main/java/com/openjiuwen/service/app/{controller/a2a/client => a2a/catalog}/RemoteAgentCatalogSnapshot.java (73%) create mode 100644 service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/RemoteAgentEntry.java rename service/agent-service-app/src/test/java/com/openjiuwen/service/app/{controller/a2a/client => a2a/catalog}/A2ARemoteAgentCardRegistryTest.java (93%) diff --git "a/documents/zh/2.\345\274\200\345\217\221\346\214\207\345\215\227/API\346\226\207\346\241\243/com.openjiuwen.service/app/a2a/catalog.README.md" "b/documents/zh/2.\345\274\200\345\217\221\346\214\207\345\215\227/API\346\226\207\346\241\243/com.openjiuwen.service/app/a2a/catalog.README.md" new file mode 100644 index 00000000..237d570b --- /dev/null +++ "b/documents/zh/2.\345\274\200\345\217\221\346\214\207\345\215\227/API\346\226\207\346\241\243/com.openjiuwen.service/app/a2a/catalog.README.md" @@ -0,0 +1,16 @@ +# a2a.catalog + +`com.openjiuwen.service.app.a2a.catalog` 提供可供 Runtime 适配器复用的远端 A2A Agent 目录。 + +## 类型 + +| Type | Description | +| --- | --- | +| `A2ARemoteAgentCardRegistry` | 线程安全地保存远端 Agent Card、调用超时与调用模式,并发布目录更新事件。 | +| `RemoteAgentEntry` | 保存单个远端 Agent 的名称、Agent Card、调用超时与调用模式。 | +| `RemoteAgentCatalogSnapshot` | 保存版本号和全量远端 Agent 条目的不可变快照。 | +| `RemoteAgentCatalogChangedEvent` | 远端 Agent 目录更新后发布的全量快照事件。 | + +## 源码路径 + +`service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/` diff --git "a/documents/zh/2.\345\274\200\345\217\221\346\214\207\345\215\227/API\346\226\207\346\241\243/com.openjiuwen.service/app/controller/a2a/client.README.md" "b/documents/zh/2.\345\274\200\345\217\221\346\214\207\345\215\227/API\346\226\207\346\241\243/com.openjiuwen.service/app/controller/a2a/client.README.md" index e0ec3651..939a2cb8 100644 --- "a/documents/zh/2.\345\274\200\345\217\221\346\214\207\345\215\227/API\346\226\207\346\241\243/com.openjiuwen.service/app/controller/a2a/client.README.md" +++ "b/documents/zh/2.\345\274\200\345\217\221\346\214\207\345\215\227/API\346\226\207\346\241\243/com.openjiuwen.service/app/controller/a2a/client.README.md" @@ -7,7 +7,6 @@ | Type | Description | | --- | --- | | `A2AAgentCardDiscovery` | 按 `openjiuwen.service.a2a.remote-agents` 拉取远端 Agent Card,失败后定时重试。 | -| `A2ARemoteAgentCardRegistry` | 保存远端 AgentCard、URL 和 timeout。 | | `A2ARemoteAgentClient` | 调用远端 Agent 的 sync / streaming client。 | ## 调用模式 diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistry.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistry.java similarity index 88% rename from service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistry.java rename to service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistry.java index 60e0efe1..bfb0191a 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistry.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistry.java @@ -2,9 +2,11 @@ * Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved. */ -package com.openjiuwen.service.app.controller.a2a.client; +package com.openjiuwen.service.app.a2a.catalog; import org.a2aproject.sdk.spec.AgentCard; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.context.ApplicationEventPublisher; import java.util.List; @@ -19,6 +21,8 @@ * @since 0.1.0 */ public class A2ARemoteAgentCardRegistry { + private static final Logger log = LoggerFactory.getLogger(A2ARemoteAgentCardRegistry.class); + /** * Default timeout in seconds for remote agent calls. */ @@ -111,12 +115,6 @@ public RemoteAgentCatalogSnapshot snapshot() { } } - /** - * A registered remote agent entry holding the card and timeout configuration. - */ - public record RemoteAgentEntry(String name, AgentCard card, int timeoutSeconds, boolean isStreaming) { - } - /** * Registers a remote agent card with a specific timeout. * @@ -146,6 +144,8 @@ public void register(String name, AgentCard card, int timeoutSeconds, boolean is } finally { updateLock.unlock(); } + log.info("Registered remote A2A Agent Card agentName={} catalogVersion={} catalogSize={} streaming={}", name, + updatedSnapshot.version(), updatedSnapshot.entries().size(), isStreaming); publishCatalogChanged(updatedSnapshot); } @@ -157,5 +157,7 @@ private RemoteAgentCatalogSnapshot createSnapshot() { private void publishCatalogChanged(RemoteAgentCatalogSnapshot updatedSnapshot) { eventPublisher.publishEvent(new RemoteAgentCatalogChangedEvent(updatedSnapshot)); + log.info("Published remote A2A Agent catalog change catalogVersion={} catalogSize={}", + updatedSnapshot.version(), updatedSnapshot.entries().size()); } } diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogChangedEvent.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/RemoteAgentCatalogChangedEvent.java similarity index 85% rename from service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogChangedEvent.java rename to service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/RemoteAgentCatalogChangedEvent.java index 04aa6abf..8616d3aa 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogChangedEvent.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/RemoteAgentCatalogChangedEvent.java @@ -2,7 +2,7 @@ * Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved. */ -package com.openjiuwen.service.app.controller.a2a.client; +package com.openjiuwen.service.app.a2a.catalog; /** * Event published after the remote A2A agent catalog changes. diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogSnapshot.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/RemoteAgentCatalogSnapshot.java similarity index 73% rename from service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogSnapshot.java rename to service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/RemoteAgentCatalogSnapshot.java index 11b22bf4..63d1799e 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCatalogSnapshot.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/RemoteAgentCatalogSnapshot.java @@ -2,7 +2,7 @@ * Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved. */ -package com.openjiuwen.service.app.controller.a2a.client; +package com.openjiuwen.service.app.a2a.catalog; import java.util.List; @@ -13,7 +13,7 @@ * @param entries complete remote-agent entries sorted by name * @since 0.1.1 */ -public record RemoteAgentCatalogSnapshot(long version, List entries) { +public record RemoteAgentCatalogSnapshot(long version, List entries) { /** * Creates an immutable snapshot. */ diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/RemoteAgentEntry.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/RemoteAgentEntry.java new file mode 100644 index 00000000..a325a8b5 --- /dev/null +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/RemoteAgentEntry.java @@ -0,0 +1,19 @@ +/* + * Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved. + */ + +package com.openjiuwen.service.app.a2a.catalog; + +import org.a2aproject.sdk.spec.AgentCard; + +/** + * Immutable remote A2A Agent registration entry. + * + * @param name remote Agent name + * @param card discovered Agent Card + * @param timeoutSeconds remote call timeout in seconds + * @param isStreaming whether Runtime should prefer streaming invocation + * @since 0.1.1 + */ +public record RemoteAgentEntry(String name, AgentCard card, int timeoutSeconds, boolean isStreaming) { +} diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java index 51519990..2ef92023 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/autoconfigure/A2AAutoConfiguration.java @@ -6,6 +6,7 @@ import com.openjiuwen.service.adapters.common.middleware.MiddlewareProperties; import com.openjiuwen.service.adapters.common.middleware.redis.RedisMiddlewareAutoConfiguration; +import com.openjiuwen.service.app.a2a.catalog.A2ARemoteAgentCardRegistry; import com.openjiuwen.service.app.config.A2AProperties; import com.openjiuwen.service.app.config.SpringEnvironmentConfigProvider; import com.openjiuwen.service.app.controller.a2a.A2AAgentExecutor; @@ -19,7 +20,6 @@ import com.openjiuwen.service.app.controller.a2a.RedisTaskStore; import com.openjiuwen.service.app.controller.a2a.WriteThrottlingTaskStore; import com.openjiuwen.service.app.controller.a2a.client.A2AAgentCardDiscovery; -import com.openjiuwen.service.app.controller.a2a.client.A2ARemoteAgentCardRegistry; import com.openjiuwen.service.app.controller.a2a.client.A2ARemoteAgentClient; import com.openjiuwen.service.app.controller.a2a.client.RemoteAgentCaller; import com.openjiuwen.service.app.controller.a2a.client.RemoteAgentCardResolver; diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2AAgentCardDiscovery.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2AAgentCardDiscovery.java index 525d7d36..6e10b60e 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2AAgentCardDiscovery.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2AAgentCardDiscovery.java @@ -4,6 +4,7 @@ package com.openjiuwen.service.app.controller.a2a.client; +import com.openjiuwen.service.app.a2a.catalog.A2ARemoteAgentCardRegistry; import com.openjiuwen.service.app.config.A2AProperties; import com.openjiuwen.service.app.config.A2AProperties.RemoteAgentProperties; import com.openjiuwen.service.spec.paths.A2AServicePaths; diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClient.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClient.java index b98c04b5..41c913ce 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClient.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClient.java @@ -4,6 +4,8 @@ package com.openjiuwen.service.app.controller.a2a.client; +import com.openjiuwen.service.app.a2a.catalog.A2ARemoteAgentCardRegistry; +import com.openjiuwen.service.app.a2a.catalog.RemoteAgentEntry; import com.openjiuwen.service.app.controller.a2a.A2aPartContent; import jakarta.annotation.PreDestroy; @@ -123,8 +125,7 @@ public A2ARemoteAgentClient(A2ARemoteAgentCardRegistry registry, int ioConcurren * @param contextId * the context/conversation ID */ - private record RemoteCallSetup(A2ARemoteAgentCardRegistry.RemoteAgentEntry entry, MessageSendParams params, - String contextId) { + private record RemoteCallSetup(RemoteAgentEntry entry, MessageSendParams params, String contextId) { } private record TaskOutcome(String taskId, TaskState state, String statusText, Task task) { @@ -189,7 +190,7 @@ private static Optional callbackConfig(RemoteCall ca * whether the client should be in streaming mode * @return the SDK client */ - private Client createClient(A2ARemoteAgentCardRegistry.RemoteAgentEntry entry, boolean isStreaming) { + private Client createClient(RemoteAgentEntry entry, boolean isStreaming) { AgentCard card = entry.card(); ClientCacheKey key = new ClientCacheKey(entry.name(), endpoint(card), isStreaming); return withApplicationClassLoader(() -> clientCache.computeIfAbsent(key, @@ -234,7 +235,7 @@ private static T withApplicationClassLoader(Supplier action) { @Override public CompletableFuture callOutcome(RemoteCall call, RemoteAgentCaller.EventObserver eventObserver) { - A2ARemoteAgentCardRegistry.RemoteAgentEntry entry = registry.get(call.agentName()) + RemoteAgentEntry entry = registry.get(call.agentName()) .orElseThrow(() -> new IllegalStateException("Unknown remote agent: " + call.agentName())); boolean isStreaming = entry.isStreaming() && call.isCallerStreaming(); return callOutcome(call, eventObserver, isStreaming); diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCardResolver.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCardResolver.java index 2c4c8388..32936701 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCardResolver.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/controller/a2a/client/RemoteAgentCardResolver.java @@ -4,6 +4,8 @@ package com.openjiuwen.service.app.controller.a2a.client; +import com.openjiuwen.service.app.a2a.catalog.A2ARemoteAgentCardRegistry; + /** * SPI for resolving a remote agent's A2A URLs by {@code agentId}. * diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistryTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistryTest.java similarity index 93% rename from service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistryTest.java rename to service/agent-service-app/src/test/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistryTest.java index 62c4e617..01710044 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentCardRegistryTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistryTest.java @@ -2,7 +2,7 @@ * Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved. */ -package com.openjiuwen.service.app.controller.a2a.client; +package com.openjiuwen.service.app.a2a.catalog; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; @@ -40,10 +40,10 @@ void registrationPublishesCompleteSortedSnapshots() { assertThat(events).hasSize(2); assertThat(events.get(0).snapshot().version()).isEqualTo(1L); - assertThat(events.get(0).snapshot().entries()).extracting(A2ARemoteAgentCardRegistry.RemoteAgentEntry::name) + assertThat(events.get(0).snapshot().entries()).extracting(RemoteAgentEntry::name) .containsExactly("weather"); assertThat(events.get(1).snapshot().version()).isEqualTo(2L); - assertThat(events.get(1).snapshot().entries()).extracting(A2ARemoteAgentCardRegistry.RemoteAgentEntry::name) + assertThat(events.get(1).snapshot().entries()).extracting(RemoteAgentEntry::name) .containsExactly("balance", "weather"); assertThat(registry.getAll()).containsExactlyElementsOf(events.get(1).snapshot().entries()); } @@ -107,7 +107,7 @@ private static A2ARemoteAgentCardRegistry registryWithEvents(List { assertThat(event.snapshot().version()).isEqualTo(1L); - assertThat(event.snapshot().entries()).extracting(A2ARemoteAgentCardRegistry.RemoteAgentEntry::name) - .containsExactly("balance"); + assertThat(event.snapshot().entries()).extracting(RemoteAgentEntry::name).containsExactly("balance"); }); }); } diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2AAgentCardDiscoveryTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2AAgentCardDiscoveryTest.java index 35326e60..634f750f 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2AAgentCardDiscoveryTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2AAgentCardDiscoveryTest.java @@ -6,6 +6,7 @@ import static org.assertj.core.api.Assertions.assertThatThrownBy; +import com.openjiuwen.service.app.a2a.catalog.A2ARemoteAgentCardRegistry; import com.openjiuwen.service.app.config.A2AProperties; import com.openjiuwen.service.app.config.A2AProperties.RemoteAgentProperties; diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientClassLoaderTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientClassLoaderTest.java index 0ba0e19c..6d38cc49 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientClassLoaderTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientClassLoaderTest.java @@ -20,6 +20,9 @@ import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import com.openjiuwen.service.app.a2a.catalog.A2ARemoteAgentCardRegistry; +import com.openjiuwen.service.app.a2a.catalog.RemoteAgentEntry; + import org.a2aproject.sdk.client.Client; import org.a2aproject.sdk.client.ClientBuilder; import org.a2aproject.sdk.client.MessageEvent; @@ -127,8 +130,7 @@ void timeoutAppliesWhileNonStreamingSdkCallIsBlocked() { var outcome = remoteClient.callOutcome(remoteCall("timeout-agent"), mock(RemoteAgentCaller.EventObserver.class)); - assertThatThrownBy(() -> outcome.get(2, TimeUnit.SECONDS)) - .hasCauseInstanceOf(TimeoutException.class); + assertThatThrownBy(() -> outcome.get(2, TimeUnit.SECONDS)).hasCauseInstanceOf(TimeoutException.class); } finally { release.countDown(); remoteClient.shutdown(); @@ -153,9 +155,8 @@ void synchronousSdkFailureCompletesOutcomeImmediately() { var outcome = remoteClient.callOutcome(remoteCall("failing-agent"), mock(RemoteAgentCaller.EventObserver.class)); - assertThatThrownBy(() -> outcome.get(1, TimeUnit.SECONDS)) - .hasCauseInstanceOf(A2AClientException.class) - .hasRootCauseMessage("SDK send failed"); + assertThatThrownBy(() -> outcome.get(1, TimeUnit.SECONDS)).hasCauseInstanceOf(A2AClientException.class) + .hasRootCauseMessage("SDK send failed"); } finally { remoteClient.shutdown(); } @@ -168,8 +169,8 @@ void synchronousSdkRuntimeFailureCompletesOutcomeImmediately() { registry.register("runtime-failing-agent", card, 30, false); ClientBuilder builder = mock(ClientBuilder.class); Client sdkClient = mock(Client.class); - doThrow(new IllegalArgumentException("invalid SDK event")) - .when(sdkClient).sendMessage(any(MessageSendParams.class), anyList(), any(), isNull()); + doThrow(new IllegalArgumentException("invalid SDK event")).when(sdkClient) + .sendMessage(any(MessageSendParams.class), anyList(), any(), isNull()); A2ARemoteAgentClient remoteClient = new A2ARemoteAgentClient(registry); try (MockedStatic clientFactory = mockStatic(Client.class)) { @@ -179,8 +180,7 @@ void synchronousSdkRuntimeFailureCompletesOutcomeImmediately() { mock(RemoteAgentCaller.EventObserver.class)); assertThatThrownBy(() -> outcome.get(1, TimeUnit.SECONDS)) - .hasCauseInstanceOf(IllegalArgumentException.class) - .hasRootCauseMessage("invalid SDK event"); + .hasCauseInstanceOf(IllegalArgumentException.class).hasRootCauseMessage("invalid SDK event"); } finally { remoteClient.shutdown(); } @@ -195,10 +195,11 @@ void directMessageEventCompletesCallWithAllTextParts() throws Exception { ClientBuilder builder = mock(ClientBuilder.class); Client sdkClient = mock(Client.class); Message message = Message.builder().role(Message.Role.ROLE_AGENT) - .parts(List.>of(new TextPart("hello "), new TextPart("world"))).build(); + .parts(List.>of(new TextPart("hello "), new TextPart("world"))).build(); doAnswer(invocation -> { - @SuppressWarnings("unchecked") List> consumers = invocation.getArgument(1); + @SuppressWarnings("unchecked") + List> consumers = invocation + .getArgument(1); consumers.get(0).accept(new MessageEvent(message), card); return null; }).when(sdkClient).sendMessage(any(MessageSendParams.class), anyList(), any(), isNull()); @@ -223,14 +224,17 @@ void completedTaskAggregatesAllArtifacts() throws Exception { ClientBuilder builder = mock(ClientBuilder.class); Client sdkClient = mock(Client.class); Task task = Task.builder().id("remote-task").contextId("remote-context") - .status(new TaskStatus(TaskState.TASK_STATE_COMPLETED)) - .artifacts(List.of( - org.a2aproject.sdk.spec.Artifact.builder().artifactId("a").parts(new TextPart("hello ")).build(), - org.a2aproject.sdk.spec.Artifact.builder().artifactId("b").parts(new TextPart("world")).build())) - .build(); + .status(new TaskStatus(TaskState.TASK_STATE_COMPLETED)) + .artifacts(List.of( + org.a2aproject.sdk.spec.Artifact.builder().artifactId("a").parts(new TextPart("hello ")) + .build(), + org.a2aproject.sdk.spec.Artifact.builder().artifactId("b").parts(new TextPart("world")) + .build())) + .build(); doAnswer(invocation -> { - @SuppressWarnings("unchecked") List> consumers = invocation.getArgument(1); + @SuppressWarnings("unchecked") + List> consumers = invocation + .getArgument(1); consumers.get(0).accept(new TaskEvent(task), card); return null; }).when(sdkClient).sendMessage(any(MessageSendParams.class), anyList(), any(), isNull()); @@ -255,14 +259,14 @@ void completedTaskWithoutArtifactsUsesStatusMessage() throws Exception { ClientBuilder builder = mock(ClientBuilder.class); Client sdkClient = mock(Client.class); Message statusMessage = Message.builder().role(Message.Role.ROLE_AGENT) - .parts(List.>of(new TextPart("status result"))).build(); + .parts(List.>of(new TextPart("status result"))).build(); Task task = Task.builder().id("remote-task").contextId("remote-context") - .status(new TaskStatus(TaskState.TASK_STATE_COMPLETED, statusMessage, null)) - .artifacts(List.of()) - .build(); + .status(new TaskStatus(TaskState.TASK_STATE_COMPLETED, statusMessage, null)).artifacts(List.of()) + .build(); doAnswer(invocation -> { - @SuppressWarnings("unchecked") List> consumers = invocation.getArgument(1); + @SuppressWarnings("unchecked") + List> consumers = invocation + .getArgument(1); consumers.get(0).accept(new TaskEvent(task), card); return null; }).when(sdkClient).sendMessage(any(MessageSendParams.class), anyList(), any(), isNull()); @@ -355,10 +359,9 @@ void createClientUsesApplicationClassLoaderForTransportDiscovery() throws Except Thread.currentThread().setContextClassLoader(new NoServicesClassLoader(original)); try { A2ARemoteAgentClient client = new A2ARemoteAgentClient(new A2ARemoteAgentCardRegistry()); - A2ARemoteAgentCardRegistry.RemoteAgentEntry entry = - new A2ARemoteAgentCardRegistry.RemoteAgentEntry("remote", testCard(), 30, true); - Method createClient = A2ARemoteAgentClient.class.getDeclaredMethod("createClient", - A2ARemoteAgentCardRegistry.RemoteAgentEntry.class, boolean.class); + RemoteAgentEntry entry = new RemoteAgentEntry("remote", testCard(), 30, true); + Method createClient = A2ARemoteAgentClient.class.getDeclaredMethod("createClient", RemoteAgentEntry.class, + boolean.class); createClient.setAccessible(true); assertThatCode(() -> createClient.invoke(client, entry, true)).doesNotThrowAnyException(); @@ -376,8 +379,8 @@ private static AgentCard testCard(String url) { .capabilities(new AgentCapabilities(true, false, false, List.of())).defaultInputModes(List.of("text")) .defaultOutputModes(List.of("text")).skills(List.of()).securitySchemes(Collections.emptyMap()) .securityRequirements(List.of()) - .supportedInterfaces(List.of(new AgentInterface("JSONRPC", url, null, "1.0"))) - .url(url).preferredTransport("JSONRPC").additionalInterfaces(List.of()).build(); + .supportedInterfaces(List.of(new AgentInterface("JSONRPC", url, null, "1.0"))).url(url) + .preferredTransport("JSONRPC").additionalInterfaces(List.of()).build(); } private static void stubClient(MockedStatic factory, AgentCard card, ClientBuilder builder, Client client) { @@ -392,8 +395,7 @@ private static RemoteCall remoteCall(String agentName) { } private static RemoteCall remoteCall(String agentName, boolean isCallerStreaming) { - return new RemoteCall(agentName, "hello", "context", null, Map.of(), Map.of(), - isCallerStreaming); + return new RemoteCall(agentName, "hello", "context", null, Map.of(), Map.of(), isCallerStreaming); } private static final class NoServicesClassLoader extends ClassLoader { diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientResultTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientResultTest.java index 8b461ce4..efaf78ca 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientResultTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientResultTest.java @@ -14,6 +14,7 @@ import com.google.gson.Gson; import com.google.gson.JsonParser; +import com.openjiuwen.service.app.a2a.catalog.A2ARemoteAgentCardRegistry; import com.openjiuwen.service.app.controller.a2a.ChunkMapper; import com.openjiuwen.service.spec.dto.QueryChunk; diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientStreamingLifecycleTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientStreamingLifecycleTest.java index de985f81..d7ec8db0 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientStreamingLifecycleTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/controller/a2a/client/A2ARemoteAgentClientStreamingLifecycleTest.java @@ -8,6 +8,7 @@ import static org.assertj.core.api.Assertions.catchThrowable; import static org.mockito.Mockito.mock; +import com.openjiuwen.service.app.a2a.catalog.A2ARemoteAgentCardRegistry; import com.sun.net.httpserver.HttpExchange; import com.sun.net.httpserver.HttpServer; diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/it/DualRuntimeCallbackIntegrationTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/it/DualRuntimeCallbackIntegrationTest.java index 15f0bc64..dcd3883c 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/it/DualRuntimeCallbackIntegrationTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/it/DualRuntimeCallbackIntegrationTest.java @@ -7,7 +7,8 @@ import static org.assertj.core.api.Assertions.assertThat; import com.fasterxml.jackson.databind.ObjectMapper; -import com.openjiuwen.service.app.controller.a2a.client.A2ARemoteAgentCardRegistry; +import com.openjiuwen.service.app.a2a.catalog.A2ARemoteAgentCardRegistry; +import com.openjiuwen.service.app.it.DualRuntimeCallbackIntegrationTest.CallerRuntimeApplication; import com.openjiuwen.service.spec.dto.QueryChunk; import com.openjiuwen.service.spec.dto.QueryResponse; import com.openjiuwen.service.spec.dto.ServeRequest; @@ -30,6 +31,7 @@ import org.springframework.boot.resttestclient.TestRestTemplate; import org.springframework.boot.resttestclient.autoconfigure.AutoConfigureTestRestTemplate; import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.context.SpringBootTest.WebEnvironment; import org.springframework.boot.test.web.server.LocalServerPort; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; @@ -52,12 +54,8 @@ /** * Dual-runtime happy path for A2A callback-mode remote invocation. */ -@SpringBootTest(classes = DualRuntimeCallbackIntegrationTest.CallerRuntimeApplication.class, - webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT, - properties = { - "spring.application.name=caller-it", - "openjiuwen.service.a2a.push-notifications=true" - }) +@SpringBootTest(classes = CallerRuntimeApplication.class, webEnvironment = WebEnvironment.RANDOM_PORT, properties = { + "spring.application.name=caller-it", "openjiuwen.service.a2a.push-notifications=true"}) @AutoConfigureTestRestTemplate @DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS) class DualRuntimeCallbackIntegrationTest { @@ -83,12 +81,8 @@ class DualRuntimeCallbackIntegrationTest { @BeforeEach void startCallee() { - callee = new SpringApplicationBuilder(CalleeRuntimeApplication.class) - .properties( - "server.port=0", - "spring.application.name=callee-it", - "openjiuwen.service.a2a.push-notifications=true") - .run(); + callee = new SpringApplicationBuilder(CalleeRuntimeApplication.class).properties("server.port=0", + "spring.application.name=callee-it", "openjiuwen.service.a2a.push-notifications=true").run(); registry.register("callee", card(calleePort()), 5, false); } @@ -103,45 +97,34 @@ void stopCallee() { @SuppressWarnings("unchecked") void callerDelegatesToCalleeWaitsForCallbackThenResumesOriginalTask() throws Exception { String callbackUrl = "http://127.0.0.1:" + callerPort + "/a2a/push-notifications/callback"; - Map firstBody = json(postA2a(rpc("SendMessage", "dual-runtime-start", Map.of( - "metadata", Map.of( - CALLBACK_URL_METADATA, callbackUrl, - CALLBACK_ID_METADATA, "push-dual-runtime"), - "message", Map.of( - "role", "ROLE_USER", - "messageId", "msg-dual-runtime-start", - "contextId", "ctx-dual-runtime", - "parts", List.of(Map.of("kind", "text", "text", "start dual runtime"))))))); + Map firstBody = json(postA2a(rpc("SendMessage", "dual-runtime-start", Map.of("metadata", + Map.of(CALLBACK_URL_METADATA, callbackUrl, CALLBACK_ID_METADATA, "push-dual-runtime"), "message", + Map.of("role", "ROLE_USER", "messageId", "msg-dual-runtime-start", "contextId", "ctx-dual-runtime", + "parts", List.of(Map.of("kind", "text", "text", "start dual runtime"))))))); Map waitingTask = taskFrom(firstBody); String taskId = String.valueOf(waitingTask.get("id")); assertThat(((Map) waitingTask.get("status")).get("state")) - .isEqualTo("TASK_STATE_INPUT_REQUIRED"); + .isEqualTo("TASK_STATE_INPUT_REQUIRED"); Map readyBatch = awaitReadyRemoteBatch(taskId); List> members = (List>) readyBatch.get("members"); assertThat(readyBatch).containsEntry("state", "READY_TO_RESUME"); - assertThat(members).singleElement().satisfies(member -> assertThat(member) - .containsEntry("agentName", "callee") - .containsEntry("state", "COMPLETED") - .containsEntry("resultCategory", "COMPLETED")); + assertThat(members).singleElement().satisfies(member -> assertThat(member).containsEntry("agentName", "callee") + .containsEntry("state", "COMPLETED").containsEntry("resultCategory", "COMPLETED")); assertThat(String.valueOf(members.get(0).get("result"))).contains("callee result:delegate:start dual runtime"); - Map resumedBody = json(postA2a(rpc("SendMessage", "dual-runtime-resume", Map.of( - "message", Map.of( - "role", "ROLE_USER", - "messageId", "msg-dual-runtime-resume", - "taskId", taskId, - "contextId", "ctx-dual-runtime", - "parts", List.of(Map.of("kind", "text", "text", "continue"))))))); + Map resumedBody = json(postA2a(rpc("SendMessage", "dual-runtime-resume", + Map.of("message", + Map.of("role", "ROLE_USER", "messageId", "msg-dual-runtime-resume", "taskId", taskId, + "contextId", "ctx-dual-runtime", "parts", + List.of(Map.of("kind", "text", "text", "continue"))))))); Map completedTask = taskFrom(resumedBody); assertThat(completedTask.get("id")).isEqualTo(taskId); - assertThat(((Map) completedTask.get("status")).get("state")) - .isEqualTo("TASK_STATE_COMPLETED"); - assertThat(allArtifactText(completedTask)) - .contains("caller resumed") - .contains("callee result:delegate:start dual runtime"); + assertThat(((Map) completedTask.get("status")).get("state")).isEqualTo("TASK_STATE_COMPLETED"); + assertThat(allArtifactText(completedTask)).contains("caller resumed") + .contains("callee result:delegate:start dual runtime"); } private Map awaitReadyRemoteBatch(String taskId) throws Exception { @@ -151,8 +134,7 @@ private Map awaitReadyRemoteBatch(String taskId) throws Exceptio while (Instant.now().isBefore(deadline)) { Task shadow = taskStore.get(shadowTaskId); if (shadow != null && shadow.metadata() != null - && shadow.metadata().get("_remote_batch") instanceof Map batch - && isCompletedBatch(batch)) { + && shadow.metadata().get("_remote_batch") instanceof Map batch && isCompletedBatch(batch)) { Map result = new LinkedHashMap<>(); batch.forEach((key, value) -> result.put(String.valueOf(key), value)); return result; @@ -160,16 +142,16 @@ && isCompletedBatch(batch)) { lastObserved = shadow == null ? "" : String.valueOf(shadow.metadata()); Thread.sleep(100); } - throw new AssertionError("remote batch was not recovered for " + shadowTaskId - + ", lastObserved=" + lastObserved); + throw new AssertionError( + "remote batch was not recovered for " + shadowTaskId + ", lastObserved=" + lastObserved); } private static boolean isCompletedBatch(Map batch) { if (!"READY_TO_RESUME".equals(batch.get("state")) || !(batch.get("members") instanceof List members)) { return false; } - return members.stream().allMatch(member -> member instanceof Map item - && "COMPLETED".equals(item.get("state"))); + return members.stream() + .allMatch(member -> member instanceof Map item && "COMPLETED".equals(item.get("state"))); } private ResponseEntity postA2a(Map body) { @@ -224,30 +206,17 @@ private static String allArtifactText(Map task) { private static AgentCard card(int port) { String url = "http://127.0.0.1:" + port + "/a2a"; - return AgentCard.builder() - .name("callee") - .description("callee") - .provider(new AgentProvider("", "")) - .version("1.0") - .capabilities(new AgentCapabilities(false, true, false, List.of())) - .defaultInputModes(List.of("text")) - .defaultOutputModes(List.of("text")) - .skills(List.of()) - .securitySchemes(Collections.emptyMap()) - .securityRequirements(List.of()) - .supportedInterfaces(List.of(new AgentInterface("JSONRPC", url, null, "1.0"))) - .url(url) - .preferredTransport("JSONRPC") - .additionalInterfaces(List.of()) - .build(); + return AgentCard.builder().name("callee").description("callee").provider(new AgentProvider("", "")) + .version("1.0").capabilities(new AgentCapabilities(false, true, false, List.of())) + .defaultInputModes(List.of("text")).defaultOutputModes(List.of("text")).skills(List.of()) + .securitySchemes(Collections.emptyMap()).securityRequirements(List.of()) + .supportedInterfaces(List.of(new AgentInterface("JSONRPC", url, null, "1.0"))).url(url) + .preferredTransport("JSONRPC").additionalInterfaces(List.of()).build(); } @SpringBootConfiguration @EnableAutoConfiguration - @ComponentScan(basePackages = { - "com.openjiuwen.service.app.controller", - "com.openjiuwen.service.app.lifecycle" - }) + @ComponentScan(basePackages = {"com.openjiuwen.service.app.controller", "com.openjiuwen.service.app.lifecycle"}) static class CallerRuntimeApplication { @Bean @Primary @@ -258,10 +227,7 @@ AgentHandler callerHandler() { @SpringBootConfiguration @EnableAutoConfiguration - @ComponentScan(basePackages = { - "com.openjiuwen.service.app.controller", - "com.openjiuwen.service.app.lifecycle" - }) + @ComponentScan(basePackages = {"com.openjiuwen.service.app.controller", "com.openjiuwen.service.app.lifecycle"}) static class CalleeRuntimeApplication { @Bean @Primary @@ -277,17 +243,13 @@ public QueryResponse query(ServeRequest request) { if (results instanceof Map remoteResults) { return response(request, "caller resumed:" + remoteResults.get("call-callee")); } - return new QueryResponse(Map.of( - "role", "assistant", - "_interrupt", Map.of( - "batchId", "dual-runtime-batch", - "items", List.of(Map.of( - "index", 0, - "toolCallId", "call-callee", - "toolName", "callee-tool", - "message", "delegate:" + request.lastUserQuery(), - "context", Map.of("_interrupt_kind", "a2a_delegate", "agentName", "callee"))))), - request.getConversationId()); + return new QueryResponse( + Map.of("role", "assistant", "_interrupt", + Map.of("batchId", "dual-runtime-batch", "items", + List.of(Map.of("index", 0, "toolCallId", "call-callee", "toolName", "callee-tool", + "message", "delegate:" + request.lastUserQuery(), "context", + Map.of("_interrupt_kind", "a2a_delegate", "agentName", "callee"))))), + request.getConversationId()); } @Override diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/it/DualRuntimeFailureIntegrationTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/it/DualRuntimeFailureIntegrationTest.java index c29dc7e0..58d21c1c 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/it/DualRuntimeFailureIntegrationTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/it/DualRuntimeFailureIntegrationTest.java @@ -7,7 +7,7 @@ import static org.assertj.core.api.Assertions.assertThat; import com.fasterxml.jackson.databind.ObjectMapper; -import com.openjiuwen.service.app.controller.a2a.client.A2ARemoteAgentCardRegistry; +import com.openjiuwen.service.app.a2a.catalog.A2ARemoteAgentCardRegistry; import com.openjiuwen.service.spec.dto.QueryChunk; import com.openjiuwen.service.spec.dto.QueryResponse; import com.openjiuwen.service.spec.dto.ServeRequest; From 7d977ab568e3ab572181ac6efb4f7b0ef4f46e39 Mon Sep 17 00:00:00 2001 From: xiaoming-2026 Date: Sat, 8 Aug 2026 11:13:03 +0800 Subject: [PATCH 5/9] revert(agentcore): keep DeepAgent invoke path --- .../adapters/agentcore/agentfw/JiuwenCoreAgentHandler.java | 7 +++---- .../agentcore/agentfw/JiuwenCoreAgentHandlerTest.java | 7 ------- 2 files changed, 3 insertions(+), 11 deletions(-) diff --git a/service/agent-service-adapters/agent-service-adapters-agentcore/src/main/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandler.java b/service/agent-service-adapters/agent-service-adapters-agentcore/src/main/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandler.java index 46a88caa..70086d01 100644 --- a/service/agent-service-adapters/agent-service-adapters-agentcore/src/main/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandler.java +++ b/service/agent-service-adapters/agent-service-adapters-agentcore/src/main/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandler.java @@ -370,11 +370,10 @@ private static QueryResponse buildQueryResponse(Object lastPayload, StringBuilde return new QueryResponse(result, conversationId); } - static boolean supportsInvoke(Object agent) { - if (agent == null || agent instanceof String || agent instanceof DeepAgent) { + private static boolean supportsInvoke(Object agent) { + if (agent == null || agent instanceof String) { // Resolved at runtime from agent-id; use streaming unless the instance exposes - // invoke. DeepAgent task-loop interruptions are emitted on its stream and may - // not be represented by the aggregate invoke result after an earlier tool call. + // invoke. return false; } for (Method method : agent.getClass().getMethods()) { diff --git a/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java b/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java index 45440680..29181ff6 100644 --- a/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java +++ b/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java @@ -26,7 +26,6 @@ import com.openjiuwen.core.singleagent.interrupt.ToolCallInterruptRequest; import com.openjiuwen.core.singleagent.schema.AgentCard; import com.openjiuwen.core.workflow.WorkflowOutput; -import com.openjiuwen.harness.deep_agent.DeepAgent; import com.openjiuwen.service.adapters.agentcore.external.ExternalSvcAdapterRegistrar; import com.openjiuwen.service.spec.dto.QueryChunk; import com.openjiuwen.service.spec.dto.QueryResponse; @@ -230,12 +229,6 @@ void syncQueryUsesInvokePathWhenAgentSupportsIt() { assertThat((Map) second.getResult()).containsEntry("content", "turn2:b|prev=a"); } - @Test - void deepAgentUsesStreamingPathToPreserveTaskLoopInterruptions() { - assertThat(JiuwenCoreAgentHandler.supportsInvoke(mock(DeepAgent.class))).isFalse(); - assertThat(JiuwenCoreAgentHandler.supportsInvoke(new InvokeEchoAgent())).isTrue(); - } - @Test @SuppressWarnings("unchecked") void nonStreamingQueryPreservesAllRemoteInterruptsInOriginalOrder() { From 197652ffedc05f2107baa5f9b077c5724713b8b6 Mon Sep 17 00:00:00 2001 From: xiaoming-2026 Date: Sat, 8 Aug 2026 11:51:27 +0800 Subject: [PATCH 6/9] fix(a2a): isolate catalog listener failures --- .../app/a2a/catalog/A2ARemoteAgentCardRegistry.java | 11 ++++++++--- .../a2a/catalog/A2ARemoteAgentCardRegistryTest.java | 13 ++++++------- 2 files changed, 14 insertions(+), 10 deletions(-) diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistry.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistry.java index bfb0191a..d7442895 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistry.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistry.java @@ -156,8 +156,13 @@ private RemoteAgentCatalogSnapshot createSnapshot() { } private void publishCatalogChanged(RemoteAgentCatalogSnapshot updatedSnapshot) { - eventPublisher.publishEvent(new RemoteAgentCatalogChangedEvent(updatedSnapshot)); - log.info("Published remote A2A Agent catalog change catalogVersion={} catalogSize={}", - updatedSnapshot.version(), updatedSnapshot.entries().size()); + try { + eventPublisher.publishEvent(new RemoteAgentCatalogChangedEvent(updatedSnapshot)); + log.info("Published remote A2A Agent catalog change catalogVersion={} catalogSize={}", + updatedSnapshot.version(), updatedSnapshot.entries().size()); + } catch (RuntimeException exception) { + log.error("Failed to publish remote A2A Agent catalog change catalogVersion={} catalogSize={}", + updatedSnapshot.version(), updatedSnapshot.entries().size(), exception); + } } } diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistryTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistryTest.java index 01710044..de6add7b 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistryTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/a2a/catalog/A2ARemoteAgentCardRegistryTest.java @@ -40,11 +40,10 @@ void registrationPublishesCompleteSortedSnapshots() { assertThat(events).hasSize(2); assertThat(events.get(0).snapshot().version()).isEqualTo(1L); - assertThat(events.get(0).snapshot().entries()).extracting(RemoteAgentEntry::name) - .containsExactly("weather"); + assertThat(events.get(0).snapshot().entries()).extracting(RemoteAgentEntry::name).containsExactly("weather"); assertThat(events.get(1).snapshot().version()).isEqualTo(2L); - assertThat(events.get(1).snapshot().entries()).extracting(RemoteAgentEntry::name) - .containsExactly("balance", "weather"); + assertThat(events.get(1).snapshot().entries()).extracting(RemoteAgentEntry::name).containsExactly("balance", + "weather"); assertThat(registry.getAll()).containsExactlyElementsOf(events.get(1).snapshot().entries()); } @@ -85,13 +84,13 @@ void concurrentRegistrationProducesUniqueCompleteVersions() { } @Test - void publicationFailureKeepsCompletedRegistryUpdateVisible() { + void publicationFailureDoesNotFailCompletedRegistryUpdate() { A2ARemoteAgentCardRegistry registry = new A2ARemoteAgentCardRegistry(event -> { throw new IllegalStateException("listener failed"); }); - assertThatThrownBy(() -> registry.register("balance", mock(AgentCard.class))) - .isInstanceOf(IllegalStateException.class).hasMessage("listener failed"); + registry.register("balance", mock(AgentCard.class)); + assertThat(registry.snapshot().version()).isEqualTo(1L); assertThat(registry.get("balance")).isPresent(); } From 9ee238a23d3bf7c18d573f2ed568f4f5fb520327 Mon Sep 17 00:00:00 2001 From: xiaoming-2026 Date: Sat, 8 Aug 2026 12:00:09 +0800 Subject: [PATCH 7/9] chore(agentcore): drop comment-only test change --- .../adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java b/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java index 29181ff6..6264656a 100644 --- a/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java +++ b/service/agent-service-adapters/agent-service-adapters-agentcore/src/test/java/com/openjiuwen/service/adapters/agentcore/agentfw/JiuwenCoreAgentHandlerTest.java @@ -44,7 +44,7 @@ import java.util.concurrent.atomic.AtomicReference; /** - * Tests AgentCore request adaptation, aggregation, and interruption handling. + * JiuwenCoreAgentHandlerTest * * @since 2026-07-03 */ From 6feb8ed3f9d8e57b5c5cc7baa196f411abb18aa0 Mon Sep 17 00:00:00 2001 From: xiaoming-2026 Date: Sat, 8 Aug 2026 12:36:38 +0800 Subject: [PATCH 8/9] feat(a2a): expose parent agent progress on interrupts --- .../orchestrator/RemoteInvocationBatch.java | 7 ++++++ .../RemoteInvocationBatchMapper.java | 17 +++++++++++--- .../RemoteInvocationBatchMapperTest.java | 22 +++++++++++++++++++ 3 files changed, 43 insertions(+), 3 deletions(-) diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatch.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatch.java index eae7b9c5..2e651ed1 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatch.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatch.java @@ -52,6 +52,8 @@ static final class Member { final String agentName; + final String parentProgress; + String message; volatile MemberState state = MemberState.QUEUED; @@ -73,11 +75,16 @@ static final class Member { Instant completedAt; Member(int index, String toolCallId, String toolName, String agentName, String message) { + this(index, toolCallId, toolName, agentName, message, ""); + } + + Member(int index, String toolCallId, String toolName, String agentName, String message, String parentProgress) { this.index = index; this.toolCallId = toolCallId; this.toolName = toolName; this.agentName = agentName; this.message = message; + this.parentProgress = parentProgress; } void fail(MemberState failedState, String category, String failureMessage) { diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapper.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapper.java index 9a0e824d..b2099847 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapper.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapper.java @@ -71,7 +71,8 @@ RemoteInvocationBatch parse(Map interrupt, ServeRequest request, } int memberIndex = item.get("index") instanceof Number number ? number.intValue() : index; members.add(new Member(memberIndex, toolCallId, stringValue(item.get("toolName")), - stringValue(context.get("agentName")), stringValue(item.get("message")))); + stringValue(context.get("agentName")), stringValue(item.get("message")), + stringValue(context.get("parentProgress")))); } members.sort(Comparator.comparingInt(member -> member.index)); return new RemoteInvocationBatch(UUID.randomUUID().toString(), parentTaskId, request, observer, members, @@ -91,7 +92,8 @@ RemoteInvocationBatch restore(Map rawBatch, ServeRequest request, String p } int index = rawMember.get("index") instanceof Number number ? number.intValue() : members.size(); Member member = new Member(index, stringValue(rawMember.get("toolCallId")), - stringValue(rawMember.get("toolName")), stringValue(rawMember.get("agentName")), ""); + stringValue(rawMember.get("toolName")), stringValue(rawMember.get("agentName")), "", + stringValue(rawMember.get("parentProgress"))); member.state = MemberState.valueOf(stringValue(rawMember.get("state"))); member.remoteTaskId = stringValue(rawMember.get("remoteTaskId")); member.resultCategory = optionalNonBlank(stringValue(rawMember.get("resultCategory"))).orElse(null); @@ -177,6 +179,7 @@ Map snapshot(RemoteInvocationBatch batch, String state) { value.put("toolCallId", member.toolCallId); value.put("toolName", member.toolName); value.put("agentName", member.agentName); + putIfNotBlank(value, "parentProgress", member.parentProgress); value.put("state", member.state.name()); putIfNotBlank(value, "remoteTaskId", member.remoteTaskId); putIfNotBlank(value, "resultCategory", member.resultCategory); @@ -249,7 +252,7 @@ private static Map publicInterrupt(RemoteInvocationBatch batch) Map item = new LinkedHashMap<>(); item.put("toolCallId", member.toolCallId); putIfNotBlank(item, "toolName", member.toolName); - item.put("message", member.inputPrompt == null ? "Remote agent requires input" : member.inputPrompt); + item.put("message", publicInputPrompt(member)); items.add(item); } Map interrupt = new LinkedHashMap<>(); @@ -259,6 +262,14 @@ private static Map publicInterrupt(RemoteInvocationBatch batch) return interrupt; } + private static String publicInputPrompt(Member member) { + String inputPrompt = member.inputPrompt == null ? "Remote agent requires input" : member.inputPrompt; + if (member.parentProgress == null || member.parentProgress.isBlank()) { + return inputPrompt; + } + return member.parentProgress + "\n\n" + inputPrompt; + } + private static Object toolResult(Member member) { if (member.state == MemberState.COMPLETED) { return member.result == null ? "" : member.result; diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapperTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapperTest.java index b56177ac..b1fdca60 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapperTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapperTest.java @@ -152,6 +152,28 @@ void resolutionMapsReadyResultsAndWaitingInterrupts() { .containsExactly(Map.of("toolCallId", "call-c", "toolName", "tool-call-c", "message", "input-c")); } + @Test + void prependsPersistedParentProgressToRemoteInputPrompt() { + Map interrupt = Map.of("toolCallId", "call-plan", "toolName", "a2a_delegate", "message", + "给张三转账100元", "context", Map.of("_interrupt_kind", "a2a_delegate", "agentName", "transfer-agent", + "parentProgress", "执行计划:\n1. 给张三转账100元\n2. 给李四转账100元")); + RemoteInvocationBatch batch = mapper.parse(interrupt, request(), "parent-plan", observer()); + Member waiting = batch.members.get(0); + waiting.state = MemberState.INPUT_REQUIRED; + waiting.inputPrompt = "请确认是否向张三转账100元。"; + + Map snapshot = mapper.snapshot(batch, "WAITING_INPUT"); + RemoteInvocationBatch restored = mapper.restore(snapshot, request(), "parent-plan", observer()); + RemoteInvocationBatchCoordinator.BatchResolution resolution = mapper.resolution(restored); + + assertThat(resolution.interrupt()).containsEntry("message", """ + 执行计划: + 1. 给张三转账100元 + 2. 给李四转账100元 + + 请确认是否向张三转账100元。"""); + } + @Test void applyOutcomeMapsRemoteStates() { Member completed = member("call-a"); From e9a91b5066a0f3c9b33915940ed05a7c84127558 Mon Sep 17 00:00:00 2001 From: xiaoming-2026 Date: Mon, 10 Aug 2026 09:24:32 +0800 Subject: [PATCH 9/9] revert(a2a): keep progress on the streaming path --- .../orchestrator/RemoteInvocationBatch.java | 7 ------ .../RemoteInvocationBatchMapper.java | 17 +++----------- .../RemoteInvocationBatchMapperTest.java | 22 ------------------- 3 files changed, 3 insertions(+), 43 deletions(-) diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatch.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatch.java index 2e651ed1..eae7b9c5 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatch.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatch.java @@ -52,8 +52,6 @@ static final class Member { final String agentName; - final String parentProgress; - String message; volatile MemberState state = MemberState.QUEUED; @@ -75,16 +73,11 @@ static final class Member { Instant completedAt; Member(int index, String toolCallId, String toolName, String agentName, String message) { - this(index, toolCallId, toolName, agentName, message, ""); - } - - Member(int index, String toolCallId, String toolName, String agentName, String message, String parentProgress) { this.index = index; this.toolCallId = toolCallId; this.toolName = toolName; this.agentName = agentName; this.message = message; - this.parentProgress = parentProgress; } void fail(MemberState failedState, String category, String failureMessage) { diff --git a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapper.java b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapper.java index b2099847..9a0e824d 100644 --- a/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapper.java +++ b/service/agent-service-app/src/main/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapper.java @@ -71,8 +71,7 @@ RemoteInvocationBatch parse(Map interrupt, ServeRequest request, } int memberIndex = item.get("index") instanceof Number number ? number.intValue() : index; members.add(new Member(memberIndex, toolCallId, stringValue(item.get("toolName")), - stringValue(context.get("agentName")), stringValue(item.get("message")), - stringValue(context.get("parentProgress")))); + stringValue(context.get("agentName")), stringValue(item.get("message")))); } members.sort(Comparator.comparingInt(member -> member.index)); return new RemoteInvocationBatch(UUID.randomUUID().toString(), parentTaskId, request, observer, members, @@ -92,8 +91,7 @@ RemoteInvocationBatch restore(Map rawBatch, ServeRequest request, String p } int index = rawMember.get("index") instanceof Number number ? number.intValue() : members.size(); Member member = new Member(index, stringValue(rawMember.get("toolCallId")), - stringValue(rawMember.get("toolName")), stringValue(rawMember.get("agentName")), "", - stringValue(rawMember.get("parentProgress"))); + stringValue(rawMember.get("toolName")), stringValue(rawMember.get("agentName")), ""); member.state = MemberState.valueOf(stringValue(rawMember.get("state"))); member.remoteTaskId = stringValue(rawMember.get("remoteTaskId")); member.resultCategory = optionalNonBlank(stringValue(rawMember.get("resultCategory"))).orElse(null); @@ -179,7 +177,6 @@ Map snapshot(RemoteInvocationBatch batch, String state) { value.put("toolCallId", member.toolCallId); value.put("toolName", member.toolName); value.put("agentName", member.agentName); - putIfNotBlank(value, "parentProgress", member.parentProgress); value.put("state", member.state.name()); putIfNotBlank(value, "remoteTaskId", member.remoteTaskId); putIfNotBlank(value, "resultCategory", member.resultCategory); @@ -252,7 +249,7 @@ private static Map publicInterrupt(RemoteInvocationBatch batch) Map item = new LinkedHashMap<>(); item.put("toolCallId", member.toolCallId); putIfNotBlank(item, "toolName", member.toolName); - item.put("message", publicInputPrompt(member)); + item.put("message", member.inputPrompt == null ? "Remote agent requires input" : member.inputPrompt); items.add(item); } Map interrupt = new LinkedHashMap<>(); @@ -262,14 +259,6 @@ private static Map publicInterrupt(RemoteInvocationBatch batch) return interrupt; } - private static String publicInputPrompt(Member member) { - String inputPrompt = member.inputPrompt == null ? "Remote agent requires input" : member.inputPrompt; - if (member.parentProgress == null || member.parentProgress.isBlank()) { - return inputPrompt; - } - return member.parentProgress + "\n\n" + inputPrompt; - } - private static Object toolResult(Member member) { if (member.state == MemberState.COMPLETED) { return member.result == null ? "" : member.result; diff --git a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapperTest.java b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapperTest.java index b1fdca60..b56177ac 100644 --- a/service/agent-service-app/src/test/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapperTest.java +++ b/service/agent-service-app/src/test/java/com/openjiuwen/service/app/orchestrator/RemoteInvocationBatchMapperTest.java @@ -152,28 +152,6 @@ void resolutionMapsReadyResultsAndWaitingInterrupts() { .containsExactly(Map.of("toolCallId", "call-c", "toolName", "tool-call-c", "message", "input-c")); } - @Test - void prependsPersistedParentProgressToRemoteInputPrompt() { - Map interrupt = Map.of("toolCallId", "call-plan", "toolName", "a2a_delegate", "message", - "给张三转账100元", "context", Map.of("_interrupt_kind", "a2a_delegate", "agentName", "transfer-agent", - "parentProgress", "执行计划:\n1. 给张三转账100元\n2. 给李四转账100元")); - RemoteInvocationBatch batch = mapper.parse(interrupt, request(), "parent-plan", observer()); - Member waiting = batch.members.get(0); - waiting.state = MemberState.INPUT_REQUIRED; - waiting.inputPrompt = "请确认是否向张三转账100元。"; - - Map snapshot = mapper.snapshot(batch, "WAITING_INPUT"); - RemoteInvocationBatch restored = mapper.restore(snapshot, request(), "parent-plan", observer()); - RemoteInvocationBatchCoordinator.BatchResolution resolution = mapper.resolution(restored); - - assertThat(resolution.interrupt()).containsEntry("message", """ - 执行计划: - 1. 给张三转账100元 - 2. 给李四转账100元 - - 请确认是否向张三转账100元。"""); - } - @Test void applyOutcomeMapsRemoteStates() { Member completed = member("call-a");