diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolAutoConfiguration.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolAutoConfiguration.java index 3250d7f308..51015334ce 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolAutoConfiguration.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolAutoConfiguration.java @@ -53,7 +53,8 @@ public class AgentProtocolAutoConfiguration { @Bean - public AgentProtocolTaskEventBus agentProtocolTaskEventBus(AgentProtocolProperties properties) { + @ConditionalOnMissingBean(AgentProtocolEventBus.class) + public AgentProtocolEventBus agentProtocolEventBus(AgentProtocolProperties properties) { return new AgentProtocolTaskEventBus(properties.getSseReplayBufferSize()); } @@ -75,7 +76,7 @@ public ProtocolTaskRepository agentProtocolTaskRepository(AgentProtocolPropertie public AgentProtocolTaskStore agentProtocolTaskStore( AgentFactory agentFactory, ProtocolTaskRepository taskRepository, - AgentProtocolTaskEventBus eventBus, + AgentProtocolEventBus eventBus, AgentProtocolProperties properties, ObjectProvider runtimeContextCustomizers) { return new AgentProtocolTaskStore( diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolEventBus.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolEventBus.java new file mode 100644 index 0000000000..6779eb04d1 --- /dev/null +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolEventBus.java @@ -0,0 +1,43 @@ +/* + * Copyright 2024-2026 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.agentscope.extensions.agentprotocol; + +import io.agentscope.harness.agent.subagent.protocol.RemoteAgentEvent; +import reactor.core.publisher.Flux; + +/** + * Event bus abstraction used by Agent Protocol task SSE endpoints. + * + *

Implementations are responsible for assigning a monotonically increasing sequence per task, + * publishing events, and replaying events after {@code fromSeq}. The default implementation is + * {@link AgentProtocolTaskEventBus}, which keeps the replay buffer in process memory. Applications + * that need cross-instance streaming can provide a shared implementation, such as a Redis Streams + * adapter, as the {@code AgentProtocolEventBus} bean. + */ +public interface AgentProtocolEventBus { + + /** Publishes an event for a task and returns the event after protocol fields are assigned. */ + RemoteAgentEvent publish(String taskId, RemoteAgentEvent event); + + /** + * Subscribes to a task's event stream, replaying only events whose sequence is greater than + * {@code fromSeq}. + */ + Flux subscribe(String taskId, long fromSeq); + + /** Completes and releases the event stream for a terminal task. */ + void complete(String taskId); +} diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskEventBus.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskEventBus.java index 4b14138a54..4dd0dcb5c7 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskEventBus.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskEventBus.java @@ -28,7 +28,7 @@ * Per-task event bus for Agent Protocol SSE streaming. Uses a replay buffer so late subscribers / * reconnects can catch up from {@code fromSeq}. */ -public final class AgentProtocolTaskEventBus { +public final class AgentProtocolTaskEventBus implements AgentProtocolEventBus { private static final Logger log = LoggerFactory.getLogger(AgentProtocolTaskEventBus.class); @@ -44,6 +44,7 @@ public AgentProtocolTaskEventBus(int replayBufferSize) { } /** Publishes an event, assigning a monotonic {@code seq}. */ + @Override public RemoteAgentEvent publish(String taskId, RemoteAgentEvent event) { Channel ch = channels.computeIfAbsent(taskId, id -> new Channel(replayBufferSize)); long seq = ch.seq.incrementAndGet(); @@ -60,6 +61,7 @@ public RemoteAgentEvent publish(String taskId, RemoteAgentEvent event) { * Subscribes to events for {@code taskId}, optionally skipping those with {@code seq <= * fromSeq}. */ + @Override public Flux subscribe(String taskId, long fromSeq) { Channel ch = channels.computeIfAbsent(taskId, id -> new Channel(replayBufferSize)); Flux flux = ch.sink.asFlux(); @@ -70,6 +72,7 @@ public Flux subscribe(String taskId, long fromSeq) { } /** Completes and removes the channel after a terminal event has been published. */ + @Override public void complete(String taskId) { Channel ch = channels.remove(taskId); if (ch != null) { diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskStore.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskStore.java index 070614dea8..481666c402 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskStore.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/main/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskStore.java @@ -73,7 +73,7 @@ public final class AgentProtocolTaskStore { private final AgentFactory agentFactory; private final ProtocolTaskRepository taskRepository; - private final AgentProtocolTaskEventBus eventBus; + private final AgentProtocolEventBus eventBus; private final AgentProtocolProperties properties; private final List runtimeContextCustomizers; private final ExecutorService executor = @@ -91,7 +91,7 @@ public final class AgentProtocolTaskStore { public AgentProtocolTaskStore( AgentFactory agentFactory, ProtocolTaskRepository taskRepository, - AgentProtocolTaskEventBus eventBus, + AgentProtocolEventBus eventBus, AgentProtocolProperties properties) { this(agentFactory, taskRepository, eventBus, properties, List.of()); } @@ -99,7 +99,7 @@ public AgentProtocolTaskStore( public AgentProtocolTaskStore( AgentFactory agentFactory, ProtocolTaskRepository taskRepository, - AgentProtocolTaskEventBus eventBus, + AgentProtocolEventBus eventBus, AgentProtocolProperties properties, List runtimeContextCustomizers) { this.agentFactory = Objects.requireNonNull(agentFactory, "agentFactory"); @@ -126,7 +126,7 @@ public AgentProtocolTaskStore( this(AgentFactory.fixed(harnessAgent), taskRepository); } - public AgentProtocolTaskEventBus eventBus() { + public AgentProtocolEventBus eventBus() { return eventBus; } diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/test/java/io/agentscope/extensions/agentprotocol/AgentProtocolAutoConfigurationTest.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/test/java/io/agentscope/extensions/agentprotocol/AgentProtocolAutoConfigurationTest.java new file mode 100644 index 0000000000..bb74cb8fed --- /dev/null +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/test/java/io/agentscope/extensions/agentprotocol/AgentProtocolAutoConfigurationTest.java @@ -0,0 +1,72 @@ +/* + * Copyright 2024-2026 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.agentscope.extensions.agentprotocol; + +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertSame; + +import io.agentscope.harness.agent.subagent.protocol.RemoteAgentEvent; +import java.util.List; +import org.junit.jupiter.api.Test; +import org.springframework.boot.autoconfigure.AutoConfigurations; +import org.springframework.boot.test.context.runner.ApplicationContextRunner; +import reactor.core.publisher.Flux; + +class AgentProtocolAutoConfigurationTest { + + private final ApplicationContextRunner contextRunner = + new ApplicationContextRunner() + .withConfiguration(AutoConfigurations.of(AgentProtocolAutoConfiguration.class)) + .withPropertyValues("agentscope.agent-protocol.enabled=true"); + + @Test + void createsInMemoryEventBusByDefault() { + contextRunner.run( + context -> + assertInstanceOf( + AgentProtocolTaskEventBus.class, + context.getBean(AgentProtocolEventBus.class))); + } + + @Test + void keepsUserEventBusWhenOneIsProvided() { + AgentProtocolEventBus customEventBus = new TestEventBus(); + + contextRunner + .withBean(AgentProtocolEventBus.class, () -> customEventBus) + .run( + context -> + assertSame( + customEventBus, + context.getBean(AgentProtocolEventBus.class))); + } + + private static final class TestEventBus implements AgentProtocolEventBus { + + @Override + public RemoteAgentEvent publish(String taskId, RemoteAgentEvent event) { + return event; + } + + @Override + public Flux subscribe(String taskId, long fromSeq) { + return Flux.fromIterable(List.of()); + } + + @Override + public void complete(String taskId) {} + } +} diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/test/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskEventBusTest.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/test/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskEventBusTest.java index 03f873d2a4..4f7e86ed5a 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/test/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskEventBusTest.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agent-protocol/src/test/java/io/agentscope/extensions/agentprotocol/AgentProtocolTaskEventBusTest.java @@ -26,6 +26,7 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.stream.LongStream; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import reactor.core.Disposable; @@ -53,6 +54,32 @@ void publish_assignsMonotonicSeqAndTaskId() { assertEquals("task-b", other.getTaskId()); } + @Test + void implementsEventBusContract() { + assertTrue(bus instanceof AgentProtocolEventBus); + } + + @Test + void publish_assignsUniqueSequencesForConcurrentPublishers() { + int eventCount = 200; + + List published = + LongStream.range(0, eventCount) + .parallel() + .mapToObj( + ignored -> + bus.publish( + "task-concurrent", event(RemoteEventType.STATUS))) + .toList(); + + assertEquals( + eventCount, published.stream().map(RemoteAgentEvent::getSeq).distinct().count()); + assertEquals( + LongStream.rangeClosed(1, eventCount).boxed().toList(), + published.stream().map(RemoteAgentEvent::getSeq).sorted().toList()); + assertTrue(published.stream().allMatch(e -> "task-concurrent".equals(e.getTaskId()))); + } + @Test void subscribe_fromSeq_skipsOlderEventsOnReplay() { bus.publish("task-replay", event(RemoteEventType.RUN_STARTED));