发布 v1.1.0 #2
@@ -257,7 +257,7 @@ public class AgentScopeReActRuntime implements AgentRuntime {
|
|||||||
AtomicReference<AgentRuntimeEvent> suspendedEvent = new AtomicReference<>();
|
AtomicReference<AgentRuntimeEvent> suspendedEvent = new AtomicReference<>();
|
||||||
// 知识库引注。
|
// 知识库引注。
|
||||||
Map<String, AgentKnowledgeReference> knowledgeReferences = new LinkedHashMap<>();
|
Map<String, AgentKnowledgeReference> knowledgeReferences = new LinkedHashMap<>();
|
||||||
// 流式输出归一化,防止出现累计快照的重复输出。
|
// 按 AgentScope 增量协议累计正文,并仅在终态快照到达时消除已发送前缀。
|
||||||
StreamDeltaNormalizer deltaNormalizer = new StreamDeltaNormalizer();
|
StreamDeltaNormalizer deltaNormalizer = new StreamDeltaNormalizer();
|
||||||
// 取消输出标记。
|
// 取消输出标记。
|
||||||
AtomicBoolean cancelled = new AtomicBoolean(false);
|
AtomicBoolean cancelled = new AtomicBoolean(false);
|
||||||
@@ -998,12 +998,12 @@ public class AgentScopeReActRuntime implements AgentRuntime {
|
|||||||
/**
|
/**
|
||||||
* 将 AgentScope 可能输出的累计快照归一化为增量。
|
* 将 AgentScope 可能输出的累计快照归一化为增量。
|
||||||
*
|
*
|
||||||
* <p>主线路 mapper 要尽量保持 AgentScope 原始顺序,但不同模型或底层适配器可能
|
* <p>当前运行时显式使用 {@code incremental(true)}。普通事件携带新增文本,必须原样保留;
|
||||||
* 输出累计文本。该归一化器只修正同一 message/block 的文本增量,不触碰旁路事件。</p>
|
* {@code last=true} 的终态事件才可能携带完整快照,此时只发送尚未输出的尾部。</p>
|
||||||
*/
|
*/
|
||||||
private static final class StreamDeltaNormalizer {
|
private static final class StreamDeltaNormalizer {
|
||||||
|
|
||||||
private final Map<String, String> previousValues = new LinkedHashMap<>();
|
private final Map<String, StringBuilder> emittedValues = new LinkedHashMap<>();
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 归一化流式事件。
|
* 归一化流式事件。
|
||||||
@@ -1030,10 +1030,14 @@ public class AgentScopeReActRuntime implements AgentRuntime {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
String key = streamKey(event, payloadKey);
|
String key = streamKey(event, payloadKey);
|
||||||
String previousText = previousValues.get(key);
|
boolean last = Boolean.TRUE.equals(event.getPayload().get("last"));
|
||||||
previousValues.put(key, currentText);
|
if (!last) {
|
||||||
if (previousText != null && !previousText.isEmpty() && currentText.startsWith(previousText)) {
|
emittedValues.computeIfAbsent(key, ignored -> new StringBuilder()).append(currentText);
|
||||||
event.getPayload().put(payloadKey, currentText.substring(previousText.length()));
|
return;
|
||||||
|
}
|
||||||
|
StringBuilder emitted = emittedValues.remove(key);
|
||||||
|
if (emitted != null && currentText.startsWith(emitted.toString())) {
|
||||||
|
event.getPayload().put(payloadKey, currentText.substring(emitted.length()));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -690,6 +690,49 @@ public class AgentScopeStatefulRuntimeTest {
|
|||||||
Assert.assertTrue(sessionStore.exists("session-1"));
|
Assert.assertTrue(sessionStore.exists("session-1"));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 验证增量模式会保留 URL 中连续出现的相同字符。
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
public void shouldPreserveRepeatedIdenticalTextDeltas() {
|
||||||
|
String expectedUrl = "http://127.0.0.1:39000/easyflow/file.docx";
|
||||||
|
AgentScopeReActRuntime runtime = runtimeWithStreamingModel(List.of(
|
||||||
|
ChatResponse.builder()
|
||||||
|
.id("url-response")
|
||||||
|
.content(List.of(TextBlock.builder().text("http://127.0.0.1:39").build()))
|
||||||
|
.build(),
|
||||||
|
ChatResponse.builder()
|
||||||
|
.id("url-response")
|
||||||
|
.content(List.of(TextBlock.builder().text("0").build()))
|
||||||
|
.build(),
|
||||||
|
ChatResponse.builder()
|
||||||
|
.id("url-response")
|
||||||
|
.content(List.of(TextBlock.builder().text("0").build()))
|
||||||
|
.build(),
|
||||||
|
ChatResponse.builder()
|
||||||
|
.id("url-response")
|
||||||
|
.content(List.of(TextBlock.builder().text("0/easyflow/file.docx").build()))
|
||||||
|
.finishReason("stop")
|
||||||
|
.build()));
|
||||||
|
runtime.init(initRequest());
|
||||||
|
|
||||||
|
List<AgentRuntimeEvent> events = runtime.stream(
|
||||||
|
AgentMessage.text(AgentMessageRole.USER, "create file"))
|
||||||
|
.collectList()
|
||||||
|
.block();
|
||||||
|
String streamedText = events.stream()
|
||||||
|
.filter(event -> event.getEventType() == AgentRuntimeEventType.MESSAGE_DELTA)
|
||||||
|
.map(event -> String.valueOf(event.getPayload().getOrDefault("text", "")))
|
||||||
|
.reduce("", String::concat);
|
||||||
|
|
||||||
|
Assert.assertEquals(expectedUrl, streamedText);
|
||||||
|
AgentRuntimeEvent completed = events.stream()
|
||||||
|
.filter(event -> event.getEventType() == AgentRuntimeEventType.COMPLETED)
|
||||||
|
.findFirst()
|
||||||
|
.orElseThrow();
|
||||||
|
Assert.assertEquals(expectedUrl, completed.getPayload().get("text"));
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void shouldAllowNextStreamAfterPreviousStreamCompleted() {
|
public void shouldAllowNextStreamAfterPreviousStreamCompleted() {
|
||||||
AgentScopeReActRuntime runtime = fakeRuntime();
|
AgentScopeReActRuntime runtime = fakeRuntime();
|
||||||
@@ -1370,6 +1413,27 @@ public class AgentScopeStatefulRuntimeTest {
|
|||||||
new AgentScopeMessageAdapter());
|
new AgentScopeMessageAdapter());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 创建单次模型调用返回多个增量响应的运行时。
|
||||||
|
*
|
||||||
|
* @param responses 同一次模型调用中的响应增量
|
||||||
|
* @return 测试运行时
|
||||||
|
*/
|
||||||
|
private AgentScopeReActRuntime runtimeWithStreamingModel(List<ChatResponse> responses) {
|
||||||
|
AgentScopeModelFactory modelFactory = new AgentScopeModelFactory() {
|
||||||
|
@Override
|
||||||
|
public Model create(AgentModelSpec modelSpec,
|
||||||
|
com.easyagents.agent.runtime.model.AgentGenerationOptions generationOptions) {
|
||||||
|
return new StreamingScriptedModel(
|
||||||
|
modelSpec == null ? "fake-model" : modelSpec.getModelName(),
|
||||||
|
responses);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
return new AgentScopeReActRuntime(modelFactory, new AgentScopeToolAdapter(),
|
||||||
|
new AgentScopeKnowledgeAdapter(), new AgentScopeMemoryAdapter(), new AgentScopeSkillAdapter(),
|
||||||
|
new AgentScopeMessageAdapter());
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 根据审批事件创建恢复请求。
|
* 根据审批事件创建恢复请求。
|
||||||
*
|
*
|
||||||
@@ -1419,6 +1483,51 @@ public class AgentScopeStatefulRuntimeTest {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 单次调用按顺序返回全部响应增量的测试模型。
|
||||||
|
*/
|
||||||
|
private static class StreamingScriptedModel implements Model {
|
||||||
|
|
||||||
|
private final String modelName;
|
||||||
|
private final List<ChatResponse> responses;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 创建流式测试模型。
|
||||||
|
*
|
||||||
|
* @param modelName 模型名称
|
||||||
|
* @param responses 响应增量
|
||||||
|
*/
|
||||||
|
private StreamingScriptedModel(String modelName, List<ChatResponse> responses) {
|
||||||
|
this.modelName = modelName;
|
||||||
|
this.responses = responses;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 返回预设的响应增量。
|
||||||
|
*
|
||||||
|
* @param messages 输入消息
|
||||||
|
* @param toolSchemas 工具定义
|
||||||
|
* @param options 生成配置
|
||||||
|
* @return 响应流
|
||||||
|
*/
|
||||||
|
@Override
|
||||||
|
public Flux<ChatResponse> stream(List<Msg> messages,
|
||||||
|
List<ToolSchema> toolSchemas,
|
||||||
|
GenerateOptions options) {
|
||||||
|
return Flux.fromIterable(responses);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 获取模型名称。
|
||||||
|
*
|
||||||
|
* @return 模型名称
|
||||||
|
*/
|
||||||
|
@Override
|
||||||
|
public String getModelName() {
|
||||||
|
return modelName;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private AgentInitRequest initRequest() {
|
private AgentInitRequest initRequest() {
|
||||||
AgentModelSpec modelSpec = new AgentModelSpec();
|
AgentModelSpec modelSpec = new AgentModelSpec();
|
||||||
modelSpec.setModelName("fake-model");
|
modelSpec.setModelName("fake-model");
|
||||||
|
|||||||
Reference in New Issue
Block a user