Skip to content

Commit 49cffeb

Browse files
hexly666jujn
andauthored
feat(agui): configure agent interruption on disconnect (#2719)
Closes #2715 --------- Co-authored-by: jujn <2087687391@qq.com>
1 parent 48b6232 commit 49cffeb

9 files changed

Lines changed: 156 additions & 30 deletions

File tree

agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/processor/AguiRequestProcessor.java

Lines changed: 24 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
*/
1616
package io.agentscope.core.agui.processor;
1717

18+
import io.agentscope.core.ReActAgent;
1819
import io.agentscope.core.agent.Agent;
1920
import io.agentscope.core.agent.RuntimeContext;
2021
import io.agentscope.core.agui.adapter.AguiAdapterConfig;
@@ -85,7 +86,29 @@ private AguiRequestProcessor(Builder builder) {
8586
* @param agent The resolved agent instance
8687
* @param events The event stream
8788
*/
88-
public record ProcessResult(Agent agent, Flux<AguiEvent> events) {}
89+
public record ProcessResult(Agent agent, Flux<AguiEvent> events) {
90+
91+
/**
92+
* Interrupt this request's active session.
93+
*
94+
* <p>AG-UI uses {@code threadId} as the session id. For a multi-session
95+
* {@link ReActAgent}, preserve the caller's user id and target that session instead of
96+
* invoking the deprecated no-argument interrupt method, which always targets the default
97+
* session.
98+
*
99+
* @param threadId The AG-UI thread id for this request
100+
* @param runtimeContext The caller-provided runtime context, may be null
101+
*/
102+
public void interrupt(String threadId, RuntimeContext runtimeContext) {
103+
if (agent instanceof ReActAgent reActAgent) {
104+
RuntimeContext interruptContext =
105+
RuntimeContext.builder(runtimeContext).sessionId(threadId).build();
106+
reActAgent.interrupt(interruptContext);
107+
} else {
108+
agent.interrupt();
109+
}
110+
}
111+
}
89112

90113
/**
91114
* Process an AG-UI request and return the result containing agent and event stream.

agentscope-extensions/agentscope-spring-boot-starters/agentscope-agui-spring-boot-starter/src/main/java/io/agentscope/spring/boot/agui/common/AguiProperties.java

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@
4040
* enable-reasoning: false
4141
* emit-token-usage: false
4242
* emit-run-finished-after-error: false
43+
* interrupt-on-disconnect: true
4344
* </pre>
4445
*/
4546
@ConfigurationProperties(prefix = "agentscope.agui")
@@ -125,6 +126,9 @@ public class AguiProperties {
125126
*/
126127
private long sseTimeout = 600000L;
127128

129+
/** Whether to interrupt the agent when the client disconnects. */
130+
private boolean interruptOnDisconnect = true;
131+
128132
public String getPathPrefix() {
129133
return pathPrefix;
130134
}
@@ -260,4 +264,12 @@ public long getSseTimeout() {
260264
public void setSseTimeout(long sseTimeout) {
261265
this.sseTimeout = sseTimeout;
262266
}
267+
268+
public boolean isInterruptOnDisconnect() {
269+
return interruptOnDisconnect;
270+
}
271+
272+
public void setInterruptOnDisconnect(boolean interruptOnDisconnect) {
273+
this.interruptOnDisconnect = interruptOnDisconnect;
274+
}
263275
}

agentscope-extensions/agentscope-spring-boot-starters/agentscope-agui-spring-boot-starter/src/main/java/io/agentscope/spring/boot/agui/mvc/AgentscopeAguiMvcAutoConfiguration.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,7 @@ public AguiMvcController aguiMvcController(
105105
.serverSideMemory(props.isServerSideMemory())
106106
.agentIdHeader(props.getAgentIdHeader())
107107
.sseTimeout(props.getSseTimeout())
108+
.interruptOnDisconnect(props.isInterruptOnDisconnect())
108109
.runtimeContextResolver(runtimeContextResolverProvider.getIfAvailable())
109110
.adapterFactory(adapterFactoryProvider.getIfAvailable())
110111
.config(config)

agentscope-extensions/agentscope-spring-boot-starters/agentscope-agui-spring-boot-starter/src/main/java/io/agentscope/spring/boot/agui/mvc/AguiMvcController.java

Lines changed: 43 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,7 @@ public class AguiMvcController {
7676
private final AguiEventEncoder encoder;
7777
private final String agentIdHeader;
7878
private final long sseTimeout;
79+
private final boolean interruptOnDisconnect;
7980
private final ExecutorService executorService;
8081
private final AguiRuntimeContextResolver runtimeContextResolver;
8182

@@ -98,6 +99,7 @@ private AguiMvcController(Builder builder) {
9899
this.agentIdHeader =
99100
builder.agentIdHeader != null ? builder.agentIdHeader : DEFAULT_AGENT_ID_HEADER;
100101
this.sseTimeout = builder.sseTimeout > 0 ? builder.sseTimeout : 600000L;
102+
this.interruptOnDisconnect = builder.interruptOnDisconnect;
101103
this.executorService = Executors.newCachedThreadPool();
102104
this.runtimeContextResolver = builder.runtimeContextResolver;
103105
}
@@ -170,34 +172,47 @@ private SseEmitter handleInternal(
170172
Disposable subscription = null;
171173
try {
172174
// Process request - returns both agent and event stream
175+
RuntimeContext runtimeContext =
176+
resolveRuntimeContext(input, headerAgentId, pathAgentId, request);
173177
AguiRequestProcessor.ProcessResult result =
174178
processor.process(
175-
input,
176-
headerAgentId,
177-
pathAgentId,
178-
resolveRuntimeContext(
179-
input, headerAgentId, pathAgentId, request));
179+
input, headerAgentId, pathAgentId, runtimeContext);
180180

181181
// Set up callbacks for client disconnect handling
182182
// using the same agent instance from the result
183183
emitter.onCompletion(
184184
() -> logger.debug("SSE connection completed for run {}", runId));
185185
emitter.onTimeout(
186186
() -> {
187-
logger.info(
188-
"SSE connection timed out for run {}, interrupting"
189-
+ " agent",
190-
runId);
191-
result.agent().interrupt();
187+
if (interruptOnDisconnect) {
188+
logger.info(
189+
"SSE connection timed out for run {}, interrupting"
190+
+ " agent",
191+
runId);
192+
result.interrupt(threadId, runtimeContext);
193+
} else {
194+
logger.info(
195+
"SSE connection timed out for run {}, agent"
196+
+ " continues running",
197+
runId);
198+
}
192199
});
193200
emitter.onError(
194201
(ex) -> {
195-
logger.info(
196-
"SSE connection error for run {}: {}, interrupting"
197-
+ " agent",
198-
runId,
199-
ex.getMessage());
200-
result.agent().interrupt();
202+
if (interruptOnDisconnect) {
203+
logger.info(
204+
"SSE connection error for run {}: {}, interrupting"
205+
+ " agent",
206+
runId,
207+
ex.getMessage());
208+
result.interrupt(threadId, runtimeContext);
209+
} else {
210+
logger.info(
211+
"SSE connection error for run {}: {}, agent"
212+
+ " continues running",
213+
runId,
214+
ex.getMessage());
215+
}
201216
});
202217

203218
// Subscribe to event stream from the same result
@@ -342,6 +357,7 @@ public static class Builder {
342357
private boolean serverSideMemory = false;
343358
private String agentIdHeader;
344359
private long sseTimeout = 600000L;
360+
private boolean interruptOnDisconnect = true;
345361
private AguiRuntimeContextResolver runtimeContextResolver;
346362
private AguiAgentAdapterFactory adapterFactory;
347363

@@ -411,6 +427,17 @@ public Builder sseTimeout(long sseTimeout) {
411427
return this;
412428
}
413429

430+
/**
431+
* Set whether to interrupt the agent when the client disconnects.
432+
*
433+
* @param interruptOnDisconnect whether to interrupt the agent
434+
* @return This builder
435+
*/
436+
public Builder interruptOnDisconnect(boolean interruptOnDisconnect) {
437+
this.interruptOnDisconnect = interruptOnDisconnect;
438+
return this;
439+
}
440+
414441
/**
415442
* Set the runtime context resolver.
416443
*

agentscope-extensions/agentscope-spring-boot-starters/agentscope-agui-spring-boot-starter/src/main/java/io/agentscope/spring/boot/agui/webflux/AgentscopeAguiWebFluxAutoConfiguration.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -107,6 +107,7 @@ public AguiWebFluxHandler aguiWebFluxHandler(
107107
.sessionManager(sessionManager)
108108
.serverSideMemory(props.isServerSideMemory())
109109
.agentIdHeader(props.getAgentIdHeader())
110+
.interruptOnDisconnect(props.isInterruptOnDisconnect())
110111
.runtimeContextResolver(runtimeContextResolverProvider.getIfAvailable())
111112
.adapterFactory(adapterFactoryProvider.getIfAvailable())
112113
.config(config)

agentscope-extensions/agentscope-spring-boot-starters/agentscope-agui-spring-boot-starter/src/main/java/io/agentscope/spring/boot/agui/webflux/AguiWebFluxHandler.java

Lines changed: 35 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -80,6 +80,7 @@ public class AguiWebFluxHandler {
8080
private final AguiRequestProcessor processor;
8181
private final AguiEventEncoder encoder;
8282
private final String agentIdHeader;
83+
private final boolean interruptOnDisconnect;
8384
private final AguiRuntimeContextResolver runtimeContextResolver;
8485

8586
private AguiWebFluxHandler(Builder builder) {
@@ -100,6 +101,7 @@ private AguiWebFluxHandler(Builder builder) {
100101
this.encoder = new AguiEventEncoder();
101102
this.agentIdHeader =
102103
builder.agentIdHeader != null ? builder.agentIdHeader : DEFAULT_AGENT_ID_HEADER;
104+
this.interruptOnDisconnect = builder.interruptOnDisconnect;
103105
this.runtimeContextResolver = builder.runtimeContextResolver;
104106
}
105107

@@ -144,31 +146,39 @@ private Mono<ServerResponse> processInput(
144146
try {
145147
// Get header agent ID
146148
String headerAgentId = request.headers().firstHeader(agentIdHeader);
149+
RuntimeContext runtimeContext =
150+
resolveRuntimeContext(input, headerAgentId, pathAgentId, request);
147151

148152
// Process request - returns both agent and event stream
149153
AguiRequestProcessor.ProcessResult result =
150-
processor.process(
151-
input,
152-
headerAgentId,
153-
pathAgentId,
154-
resolveRuntimeContext(input, headerAgentId, pathAgentId, request));
154+
processor.process(input, headerAgentId, pathAgentId, runtimeContext);
155155

156156
// Create SSE stream using ServerSentEvent for proper streaming behavior
157+
Flux<AguiEvent> events =
158+
interruptOnDisconnect
159+
? result.events()
160+
: result.events().publish().autoConnect(1);
157161
Flux<ServerSentEvent<String>> sseStream =
158-
result.events()
159-
.map(
162+
events.map(
160163
event ->
161164
ServerSentEvent.<String>builder()
162165
.data(encoder.encodeToJson(event).trim())
163166
.build())
164-
// When client closes connection (cancels stream), interrupt the agent
167+
// When the client closes the connection, optionally interrupt the agent
165168
.doOnCancel(
166169
() -> {
167-
logger.info(
168-
"SSE stream cancelled for run {}, interrupting"
169-
+ " agent",
170-
runId);
171-
result.agent().interrupt();
170+
if (interruptOnDisconnect) {
171+
logger.info(
172+
"SSE stream cancelled for run {}, interrupting"
173+
+ " agent",
174+
runId);
175+
result.interrupt(threadId, runtimeContext);
176+
} else {
177+
logger.info(
178+
"SSE stream cancelled for run {}, agent"
179+
+ " continues running",
180+
runId);
181+
}
172182
});
173183

174184
return ServerResponse.ok()
@@ -284,6 +294,7 @@ public static class Builder {
284294
private AguiAdapterConfig config;
285295
private boolean serverSideMemory = false;
286296
private String agentIdHeader;
297+
private boolean interruptOnDisconnect = true;
287298
private AguiRuntimeContextResolver runtimeContextResolver;
288299
private AguiAgentAdapterFactory adapterFactory;
289300

@@ -342,6 +353,17 @@ public Builder agentIdHeader(String agentIdHeader) {
342353
return this;
343354
}
344355

356+
/**
357+
* Set whether to interrupt the agent when the client disconnects.
358+
*
359+
* @param interruptOnDisconnect whether to interrupt the agent
360+
* @return This builder
361+
*/
362+
public Builder interruptOnDisconnect(boolean interruptOnDisconnect) {
363+
this.interruptOnDisconnect = interruptOnDisconnect;
364+
return this;
365+
}
366+
345367
/**
346368
* Set the runtime context resolver.
347369
*

agentscope-extensions/agentscope-spring-boot-starters/agentscope-agui-spring-boot-starter/src/test/java/io/agentscope/spring/boot/agui/common/AguiAdapterConfigAutoConfigurationTest.java

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,9 +85,34 @@ void testDefaultAdapterConfig() {
8585
assertFalse(config.isBaseEventPropertiesEnricherEnabled());
8686
assertTrue(config.getEventConverters().isEmpty());
8787
assertTrue(config.getEventEnrichers().isEmpty());
88+
assertTrue(interruptOnDisconnect(ctx.getBean(AguiMvcController.class)));
8889
});
8990
}
9091

92+
@Test
93+
@DisplayName("Should bind interrupt-on-disconnect property for MVC")
94+
void testMvcInterruptOnDisconnectProperty() {
95+
mvcContextRunner
96+
.withPropertyValues("agentscope.agui.interrupt-on-disconnect=false")
97+
.run(
98+
ctx ->
99+
assertFalse(
100+
interruptOnDisconnect(
101+
ctx.getBean(AguiMvcController.class))));
102+
}
103+
104+
@Test
105+
@DisplayName("Should bind interrupt-on-disconnect property for WebFlux")
106+
void testWebFluxInterruptOnDisconnectProperty() {
107+
webFluxContextRunner
108+
.withPropertyValues("agentscope.agui.interrupt-on-disconnect=false")
109+
.run(
110+
ctx ->
111+
assertFalse(
112+
interruptOnDisconnect(
113+
ctx.getBean(AguiWebFluxHandler.class))));
114+
}
115+
91116
@Test
92117
@DisplayName("Should bind emit-token-usage property for MVC")
93118
void testMvcEmitTokenUsageProperty() {
@@ -376,6 +401,10 @@ private static AguiRuntimeContextResolver webFluxRuntimeContextResolver(
376401
ReflectionTestUtils.getField(handler, "runtimeContextResolver");
377402
}
378403

404+
private static boolean interruptOnDisconnect(Object handler) {
405+
return (boolean) ReflectionTestUtils.getField(handler, "interruptOnDisconnect");
406+
}
407+
379408
private static AguiAgentAdapterFactory mvcAdapterFactory(AguiMvcController controller) {
380409
Object processor = ReflectionTestUtils.getField(controller, "processor");
381410
return (AguiAgentAdapterFactory) ReflectionTestUtils.getField(processor, "adapterFactory");

docs/v2/en/integration/protocol/agui.md

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -220,8 +220,14 @@ agentscope:
220220
enable-reasoning: false
221221
emit-run-finished-after-error: false
222222
server-side-memory: false
223+
interrupt-on-disconnect: true
223224
```
224225

226+
`interrupt-on-disconnect` controls whether an Agent run is interrupted when the MVC/WebFlux SSE
227+
connection is closed, times out, or fails while sending an event. It defaults to `true` for
228+
backward compatibility. Set it to `false` to let the Agent continue running after the client
229+
disconnects; events produced while the connection is closed are not replayed by the starter.
230+
225231
You can extend the default chain with beans:
226232

227233
- `AgentEventConverter`: custom event semantic mapping.

docs/v2/zh/integration/protocol/agui.md

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -220,8 +220,13 @@ agentscope:
220220
enable-reasoning: false
221221
emit-run-finished-after-error: false
222222
server-side-memory: false
223+
interrupt-on-disconnect: true
223224
```
224225

226+
`interrupt-on-disconnect` 用于控制 MVC/WebFlux 的 SSE 连接关闭、超时或发送事件失败时是否中断
227+
Agent run。默认值为 `true`,用于保持现有行为兼容。设置为 `false` 后,客户端断开时 Agent
228+
会继续执行;连接关闭期间产生的事件不会由 starter 重放。
229+
225230
可以通过 bean 扩展默认链路:
226231

227232
- `AgentEventConverter`:注册自定义事件语义映射。

0 commit comments

Comments
 (0)