From c6290784c0bcecaf145fc7fe3a22d28ef56bddba Mon Sep 17 00:00:00 2001 From: Zhang ChuanJin Date: Sun, 6 Sep 2026 08:22:38 +0800 Subject: [PATCH 1/3] fix(agui): use unique ids for text segments --- .../java/io/agentscope/core/ReActAgent.java | 22 +++++-- .../agent/ReActAgentNewLoopReplyTest.java | 57 +++++++++++++++++++ .../adapter/strategy/AguiStreamContext.java | 24 ++++++++ .../strategy/TextBlockEventConverter.java | 6 +- .../agui/adapter/AguiAgentAdapterV2Test.java | 54 ++++++++++++++++++ 5 files changed, 157 insertions(+), 6 deletions(-) diff --git a/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java b/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java index 01143e4207..3fa1aefd44 100644 --- a/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java +++ b/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java @@ -2593,7 +2593,10 @@ private void emitBlockEvents( blockLifecycle.startText(events); if (tb.getText() != null && !tb.getText().isEmpty()) { events.add( - new TextBlockDeltaEvent(blockLifecycle.replyId, "text", tb.getText())); + new TextBlockDeltaEvent( + blockLifecycle.replyId, + blockLifecycle.currentTextBlockId(), + tb.getText())); } } else if (block instanceof ThinkingBlock tb) { blockLifecycle.startThinking(events); @@ -2628,6 +2631,8 @@ private void emitBlockEvents( private final class ModelCallBlockLifecycle { private final String replyId; private final AtomicBoolean textStarted = new AtomicBoolean(false); + private final AtomicLong textSegmentSequence = new AtomicLong(0); + private final AtomicReference currentTextBlockId = new AtomicReference<>(); private final AtomicBoolean thinkingStarted = new AtomicBoolean(false); private final Map startedToolCalls = new ConcurrentHashMap<>(); @@ -2638,10 +2643,17 @@ private ModelCallBlockLifecycle(String replyId) { private void startText(List events) { flushThinking(events); if (textStarted.compareAndSet(false, true)) { - events.add(new TextBlockStartEvent(replyId, "text")); + long segment = textSegmentSequence.incrementAndGet(); + String blockId = segment == 1 ? "text" : "text-" + segment; + currentTextBlockId.set(blockId); + events.add(new TextBlockStartEvent(replyId, blockId)); } } + private String currentTextBlockId() { + return currentTextBlockId.get(); + } + private void startThinking(List events) { if (thinkingStarted.compareAndSet(false, true)) { events.add(new ThinkingBlockStartEvent(replyId, "thinking")); @@ -2663,7 +2675,8 @@ private void startToolCall(String toolId, String toolName, List even private void flushText(List events) { if (textStarted.compareAndSet(true, false)) { - events.add(new TextBlockEndEvent(replyId, "text")); + String blockId = currentTextBlockId.getAndSet(null); + events.add(new TextBlockEndEvent(replyId, blockId)); } } @@ -3652,7 +3665,8 @@ private Flux summaryModelCallStream( new TextBlockDeltaEvent( blockLifecycle .replyId, - "text", + blockLifecycle + .currentTextBlockId(), tb.getText())); } } else if (block diff --git a/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopReplyTest.java b/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopReplyTest.java index 5179068c0f..9df16a5a68 100644 --- a/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopReplyTest.java +++ b/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopReplyTest.java @@ -29,6 +29,7 @@ import io.agentscope.core.event.ModelCallEndEvent; import io.agentscope.core.event.ModelCallStartEvent; import io.agentscope.core.event.RequireExternalExecutionEvent; +import io.agentscope.core.event.TextBlockDeltaEvent; import io.agentscope.core.event.TextBlockEndEvent; import io.agentscope.core.event.TextBlockStartEvent; import io.agentscope.core.event.ThinkingBlockEndEvent; @@ -517,6 +518,62 @@ void consecutiveTextChunksEmitSingleStartAndEnd() { < indexOf(events, ModelCallEndEvent.class)); } + @Test + void textSeparatedByToolCallUsesDistinctBlockIds() { + ChatModelBase model = + new ScriptedModel( + List.of( + () -> + Flux.just( + chatResponse( + TextBlock.builder().text("before").build()), + chatResponse( + ToolUseBlock.builder() + .id("tc1") + .name("echo") + .input(Map.of("query", "ping")) + .build()), + chatResponse( + TextBlock.builder().text("after").build())), + () -> Flux.just(textResponse("done")))); + ReActAgent agent = + ReActAgent.builder() + .name("asst") + .model(model) + .toolkit(toolkitWith(new EchoTool())) + .build(); + + List events = agent.streamEvents(List.of()).collectList().block(); + assertNotNull(events); + + int firstModelEnd = indexOf(events, ModelCallEndEvent.class); + List starts = + events.subList(0, firstModelEnd).stream() + .filter(TextBlockStartEvent.class::isInstance) + .map(TextBlockStartEvent.class::cast) + .toList(); + List ends = + events.subList(0, firstModelEnd).stream() + .filter(TextBlockEndEvent.class::isInstance) + .map(TextBlockEndEvent.class::cast) + .toList(); + List deltas = + events.subList(0, firstModelEnd).stream() + .filter(TextBlockDeltaEvent.class::isInstance) + .map(TextBlockDeltaEvent.class::cast) + .toList(); + + assertEquals( + List.of("text", "text-2"), + starts.stream().map(TextBlockStartEvent::getBlockId).toList()); + assertEquals( + List.of("text", "text-2"), + deltas.stream().map(TextBlockDeltaEvent::getBlockId).toList()); + assertEquals( + List.of("text", "text-2"), + ends.stream().map(TextBlockEndEvent::getBlockId).toList()); + } + @Test void summaryModelCallClosesThinkingBeforeTextAndFlushesTextBeforeModelEnd() { ChatModelBase model = diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java index 11c041b279..d7510e5650 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java @@ -53,6 +53,8 @@ public class AguiStreamContext { private final Set startedTextMessages = new LinkedHashSet<>(); private final Set endedTextMessages = new LinkedHashSet<>(); + private final Map textMessageIds = new LinkedHashMap<>(); + private final Map firstTextBlockByReply = new LinkedHashMap<>(); private final Set startedReasoningMessages = new LinkedHashSet<>(); private final Set endedReasoningMessages = new LinkedHashSet<>(); private final Set startedToolCalls = new LinkedHashSet<>(); @@ -148,6 +150,26 @@ public void appendTextDelta(String messageId, String delta) { } } + public String textMessageId(String replyId, String blockId) { + String normalizedBlockId = isBlank(blockId) ? "text" : blockId; + TextBlockKey key = new TextBlockKey(replyId, normalizedBlockId); + return textMessageIds.computeIfAbsent( + key, + ignored -> { + String firstBlockId = + firstTextBlockByReply.putIfAbsent(replyId, normalizedBlockId); + if (firstBlockId == null || Objects.equals(firstBlockId, normalizedBlockId)) { + return replyId; + } + return replyId + "-" + normalizedBlockId; + }); + } + + public String existingTextMessageId(String replyId, String blockId) { + String normalizedBlockId = isBlank(blockId) ? "text" : blockId; + return textMessageIds.get(new TextBlockKey(replyId, normalizedBlockId)); + } + public void closeActiveTextMessage() { if (currentTextMessageId == null) { return; @@ -357,6 +379,8 @@ private static boolean isBlank(String value) { return value == null || value.isBlank(); } + private record TextBlockKey(String replyId, String blockId) {} + private void warnMissingToolCallId(String eventName) { if (!warnedMissingToolCallIdOperations.add(eventName)) { return; diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/TextBlockEventConverter.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/TextBlockEventConverter.java index 438c758c1f..6980b4343f 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/TextBlockEventConverter.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/TextBlockEventConverter.java @@ -33,9 +33,11 @@ public Set> eventTypes() { public void convert(AgentEvent event, AguiStreamContext context) { if (event instanceof TextBlockDeltaEvent delta) { // AguiEvent.TextMessageStart delays sending when content arrives - context.appendTextDelta(delta.getReplyId(), delta.getDelta()); + String messageId = context.textMessageId(delta.getReplyId(), delta.getBlockId()); + context.appendTextDelta(messageId, delta.getDelta()); } else if (event instanceof TextBlockEndEvent end) { - context.closeTextMessage(end.getReplyId()); + String messageId = context.existingTextMessageId(end.getReplyId(), end.getBlockId()); + context.closeTextMessage(messageId); } } } diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java index 62d5a1dfe2..a4e6af6dbf 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java @@ -329,6 +329,60 @@ void testTextBlockEventsConvertToAguiTextMessageEvents() { assertEquals("hel", firstDelta.delta()); } + @Test + void testTextSegmentsSeparatedByToolCallUseDistinctMessageIds() { + List events = + runReActEvents( + new TextBlockStartEvent("reply-mixed", "text"), + new TextBlockDeltaEvent("reply-mixed", "text", "before"), + new TextBlockEndEvent("reply-mixed", "text"), + new ToolCallStartEvent("reply-mixed", "tool-1", "lookup"), + new ToolCallEndEvent("reply-mixed", "tool-1", "lookup"), + new TextBlockStartEvent("reply-mixed", "text-2"), + new TextBlockDeltaEvent("reply-mixed", "text-2", "after"), + new TextBlockEndEvent("reply-mixed", "text-2")); + + assertEquals( + List.of( + AguiEventType.TEXT_MESSAGE_START, + AguiEventType.TEXT_MESSAGE_CONTENT, + AguiEventType.TEXT_MESSAGE_END, + AguiEventType.TOOL_CALL_START, + AguiEventType.TOOL_CALL_END, + AguiEventType.TEXT_MESSAGE_START, + AguiEventType.TEXT_MESSAGE_CONTENT, + AguiEventType.TEXT_MESSAGE_END), + types(events)); + + List messageIds = + events.stream() + .filter( + event -> + event instanceof AguiEvent.TextMessageStart + || event instanceof AguiEvent.TextMessageContent + || event instanceof AguiEvent.TextMessageEnd) + .map( + event -> { + if (event instanceof AguiEvent.TextMessageStart start) { + return start.messageId(); + } + if (event instanceof AguiEvent.TextMessageContent content) { + return content.messageId(); + } + return ((AguiEvent.TextMessageEnd) event).messageId(); + }) + .toList(); + assertEquals( + List.of( + "reply-mixed", + "reply-mixed", + "reply-mixed", + "reply-mixed-text-2", + "reply-mixed-text-2", + "reply-mixed-text-2"), + messageIds); + } + @Test void testThinkingEventsAreIgnoredWhenReasoningDisabled() { List events = From 31bbf1ef0a3be2f1d2525246d2841486b2d1d5eb Mon Sep 17 00:00:00 2001 From: Zhang ChuanJin Date: Sun, 6 Sep 2026 19:56:00 +0800 Subject: [PATCH 2/3] fix(agui): use unique ids for reasoning segments --- .../java/io/agentscope/core/ReActAgent.java | 25 ++++++-- .../agent/ReActAgentNewLoopReplyTest.java | 61 ++++++++++++++++++ .../adapter/strategy/AguiStreamContext.java | 41 ++++++++++-- .../strategy/ThinkingBlockEventConverter.java | 7 ++- .../agui/adapter/AguiAgentAdapterV2Test.java | 62 +++++++++++++++++++ 5 files changed, 184 insertions(+), 12 deletions(-) diff --git a/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java b/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java index 3fa1aefd44..406ab08877 100644 --- a/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java +++ b/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java @@ -2603,7 +2603,9 @@ private void emitBlockEvents( if (tb.getThinking() != null && !tb.getThinking().isEmpty()) { events.add( new ThinkingBlockDeltaEvent( - blockLifecycle.replyId, "thinking", tb.getThinking())); + blockLifecycle.replyId, + blockLifecycle.currentThinkingBlockId(), + tb.getThinking())); } } else if (withToolEvents && block instanceof ToolUseBlock tub) { String toolId = resolveToolCallId(tub, context); @@ -2625,8 +2627,8 @@ private void emitBlockEvents( * *

The model stream is consumed through {@code concatMap}, but the state holders keep the * previous thread-safe shape because model providers may deliver chunk content - * unpredictably. This helper only changes when pending end events are flushed; it does not - * change the block identity or event payloads. + * unpredictably. Each contiguous text or thinking segment receives its own block ID so its + * start, delta, and end events can be correlated independently. */ private final class ModelCallBlockLifecycle { private final String replyId; @@ -2634,6 +2636,8 @@ private final class ModelCallBlockLifecycle { private final AtomicLong textSegmentSequence = new AtomicLong(0); private final AtomicReference currentTextBlockId = new AtomicReference<>(); private final AtomicBoolean thinkingStarted = new AtomicBoolean(false); + private final AtomicLong thinkingSegmentSequence = new AtomicLong(0); + private final AtomicReference currentThinkingBlockId = new AtomicReference<>(); private final Map startedToolCalls = new ConcurrentHashMap<>(); private ModelCallBlockLifecycle(String replyId) { @@ -2656,10 +2660,17 @@ private String currentTextBlockId() { private void startThinking(List events) { if (thinkingStarted.compareAndSet(false, true)) { - events.add(new ThinkingBlockStartEvent(replyId, "thinking")); + long segment = thinkingSegmentSequence.incrementAndGet(); + String blockId = segment == 1 ? "thinking" : "thinking-" + segment; + currentThinkingBlockId.set(blockId); + events.add(new ThinkingBlockStartEvent(replyId, blockId)); } } + private String currentThinkingBlockId() { + return currentThinkingBlockId.get(); + } + private void startToolCall(String toolId, String toolName, List events) { if (toolId == null || startedToolCalls.containsKey(toolId)) { return; @@ -2682,7 +2693,8 @@ private void flushText(List events) { private void flushThinking(List events) { if (thinkingStarted.compareAndSet(true, false)) { - events.add(new ThinkingBlockEndEvent(replyId, "thinking")); + String blockId = currentThinkingBlockId.getAndSet(null); + events.add(new ThinkingBlockEndEvent(replyId, blockId)); } } @@ -3680,7 +3692,8 @@ private Flux summaryModelCallStream( new ThinkingBlockDeltaEvent( blockLifecycle .replyId, - "thinking", + blockLifecycle + .currentThinkingBlockId(), tb .getThinking())); } diff --git a/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopReplyTest.java b/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopReplyTest.java index 9df16a5a68..01d69d9a57 100644 --- a/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopReplyTest.java +++ b/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopReplyTest.java @@ -32,6 +32,7 @@ import io.agentscope.core.event.TextBlockDeltaEvent; import io.agentscope.core.event.TextBlockEndEvent; import io.agentscope.core.event.TextBlockStartEvent; +import io.agentscope.core.event.ThinkingBlockDeltaEvent; import io.agentscope.core.event.ThinkingBlockEndEvent; import io.agentscope.core.event.ThinkingBlockStartEvent; import io.agentscope.core.event.ToolCallEndEvent; @@ -574,6 +575,66 @@ void textSeparatedByToolCallUsesDistinctBlockIds() { ends.stream().map(TextBlockEndEvent::getBlockId).toList()); } + @Test + void thinkingSeparatedByToolCallUsesDistinctBlockIds() { + ChatModelBase model = + new ScriptedModel( + List.of( + () -> + Flux.just( + chatResponse( + ThinkingBlock.builder() + .thinking("before") + .build()), + chatResponse( + ToolUseBlock.builder() + .id("tc1") + .name("echo") + .input(Map.of("query", "ping")) + .build()), + chatResponse( + ThinkingBlock.builder() + .thinking("after") + .build())), + () -> Flux.just(textResponse("done")))); + ReActAgent agent = + ReActAgent.builder() + .name("asst") + .model(model) + .toolkit(toolkitWith(new EchoTool())) + .build(); + + List events = agent.streamEvents(List.of()).collectList().block(); + assertNotNull(events); + + int firstModelEnd = indexOf(events, ModelCallEndEvent.class); + List starts = + events.subList(0, firstModelEnd).stream() + .filter(ThinkingBlockStartEvent.class::isInstance) + .map(ThinkingBlockStartEvent.class::cast) + .toList(); + List ends = + events.subList(0, firstModelEnd).stream() + .filter(ThinkingBlockEndEvent.class::isInstance) + .map(ThinkingBlockEndEvent.class::cast) + .toList(); + List deltas = + events.subList(0, firstModelEnd).stream() + .filter(ThinkingBlockDeltaEvent.class::isInstance) + .map(ThinkingBlockDeltaEvent.class::cast) + .toList(); + + assertEquals( + List.of("thinking", "thinking-2"), + starts.stream().map(ThinkingBlockStartEvent::getBlockId).toList()); + assertEquals( + List.of("thinking", "thinking-2"), + deltas.stream().map(ThinkingBlockDeltaEvent::getBlockId).toList()); + assertEquals( + List.of("thinking", "thinking-2"), + ends.stream().map(ThinkingBlockEndEvent::getBlockId).toList()); + } + @Test void summaryModelCallClosesThinkingBeforeTextAndFlushesTextBeforeModelEnd() { ChatModelBase model = diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java index d7510e5650..039e7b0eec 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java @@ -57,6 +57,8 @@ public class AguiStreamContext { private final Map firstTextBlockByReply = new LinkedHashMap<>(); private final Set startedReasoningMessages = new LinkedHashSet<>(); private final Set endedReasoningMessages = new LinkedHashSet<>(); + private final Map reasoningMessageIds = new LinkedHashMap<>(); + private final Map firstThinkingBlockByReply = new LinkedHashMap<>(); private final Set startedToolCalls = new LinkedHashSet<>(); private final Set endedToolCalls = new LinkedHashSet<>(); private String currentTextMessageId; @@ -191,7 +193,7 @@ public void closeTextMessage(String messageId) { } public void startReasoningMessage(String messageId) { - String reasoningMessageId = reasoningMessageId(messageId); + String reasoningMessageId = withReasoningSuffix(messageId); if (startedReasoningMessages.add(reasoningMessageId)) { emit( new AguiEvent.ReasoningMessageStart( @@ -205,10 +207,36 @@ public void appendReasoningDelta(String messageId, String delta) { startReasoningMessage(messageId); emit( new AguiEvent.ReasoningMessageContent( - threadId, runId, reasoningMessageId(messageId), delta)); + threadId, runId, withReasoningSuffix(messageId), delta)); } } + public String reasoningMessageId(String replyId, String blockId) { + String normalizedBlockId = isBlank(blockId) ? "thinking" : blockId; + ThinkingBlockKey key = new ThinkingBlockKey(replyId, normalizedBlockId); + String messageId = + reasoningMessageIds.computeIfAbsent( + key, + ignored -> { + String firstBlockId = + firstThinkingBlockByReply.putIfAbsent( + replyId, normalizedBlockId); + if (firstBlockId == null + || Objects.equals(firstBlockId, normalizedBlockId)) { + return replyId; + } + return replyId + "-" + normalizedBlockId; + }); + return withReasoningSuffix(messageId); + } + + public String existingReasoningMessageId(String replyId, String blockId) { + String normalizedBlockId = isBlank(blockId) ? "thinking" : blockId; + String messageId = + reasoningMessageIds.get(new ThinkingBlockKey(replyId, normalizedBlockId)); + return messageId == null ? null : withReasoningSuffix(messageId); + } + public void closeActiveReasoningMessage() { if (currentReasoningMessageId == null) { return; @@ -217,7 +245,7 @@ public void closeActiveReasoningMessage() { } public void closeReasoningMessage(String messageId) { - String reasoningMessageId = reasoningMessageId(messageId); + String reasoningMessageId = withReasoningSuffix(messageId); if (reasoningMessageId == null || !startedReasoningMessages.contains(reasoningMessageId) || endedReasoningMessages.contains(reasoningMessageId)) { @@ -357,7 +385,10 @@ private static String normalizeToolCallName(String toolCallName) { return toolCallName != null && !toolCallName.isBlank() ? toolCallName : "unknown"; } - private static String reasoningMessageId(String messageId) { + private static String withReasoningSuffix(String messageId) { + if (messageId == null) { + return null; + } if (messageId.endsWith(REASONING_MESSAGE_ID_SUFFIX)) { return messageId; } @@ -381,6 +412,8 @@ private static boolean isBlank(String value) { private record TextBlockKey(String replyId, String blockId) {} + private record ThinkingBlockKey(String replyId, String blockId) {} + private void warnMissingToolCallId(String eventName) { if (!warnedMissingToolCallIdOperations.add(eventName)) { return; diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ThinkingBlockEventConverter.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ThinkingBlockEventConverter.java index bdf7ef4132..12e08b2bb6 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ThinkingBlockEventConverter.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ThinkingBlockEventConverter.java @@ -40,9 +40,12 @@ public void convert(AgentEvent event, AguiStreamContext context) { if (event instanceof ThinkingBlockDeltaEvent delta) { // AguiEvent.ReasoningMessageStart delays sending when content arrives - context.appendReasoningDelta(delta.getReplyId(), delta.getDelta()); + String messageId = context.reasoningMessageId(delta.getReplyId(), delta.getBlockId()); + context.appendReasoningDelta(messageId, delta.getDelta()); } else if (event instanceof ThinkingBlockEndEvent end) { - context.closeReasoningMessage(end.getReplyId()); + String messageId = + context.existingReasoningMessageId(end.getReplyId(), end.getBlockId()); + context.closeReasoningMessage(messageId); } } } diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java index a4e6af6dbf..c8bf95b823 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java @@ -383,6 +383,68 @@ void testTextSegmentsSeparatedByToolCallUseDistinctMessageIds() { messageIds); } + @Test + void testReasoningSegmentsSeparatedByToolCallUseDistinctMessageIds() { + List events = + runReActEvents( + AguiAdapterConfig.builder().enableReasoning(true).build(), + new ThinkingBlockStartEvent("reply-mixed", "thinking"), + new ThinkingBlockDeltaEvent("reply-mixed", "thinking", "before"), + new ThinkingBlockEndEvent("reply-mixed", "thinking"), + new ToolCallStartEvent("reply-mixed", "tool-1", "lookup"), + new ToolCallEndEvent("reply-mixed", "tool-1", "lookup"), + new ThinkingBlockStartEvent("reply-mixed", "thinking-2"), + new ThinkingBlockDeltaEvent("reply-mixed", "thinking-2", "after"), + new ThinkingBlockEndEvent("reply-mixed", "thinking-2")); + + assertEquals( + List.of( + AguiEventType.REASONING_MESSAGE_START, + AguiEventType.REASONING_MESSAGE_CONTENT, + AguiEventType.REASONING_MESSAGE_END, + AguiEventType.TOOL_CALL_START, + AguiEventType.TOOL_CALL_END, + AguiEventType.REASONING_MESSAGE_START, + AguiEventType.REASONING_MESSAGE_CONTENT, + AguiEventType.REASONING_MESSAGE_END), + types(events)); + + List messageIds = + events.stream() + .filter( + event -> + event instanceof AguiEvent.ReasoningMessageStart + || event + instanceof + AguiEvent.ReasoningMessageContent + || event + instanceof + AguiEvent.ReasoningMessageEnd) + .map( + event -> { + if (event + instanceof AguiEvent.ReasoningMessageStart start) { + return start.messageId(); + } + if (event + instanceof + AguiEvent.ReasoningMessageContent content) { + return content.messageId(); + } + return ((AguiEvent.ReasoningMessageEnd) event).messageId(); + }) + .toList(); + assertEquals( + List.of( + "reply-mixed-reasoning", + "reply-mixed-reasoning", + "reply-mixed-reasoning", + "reply-mixed-thinking-2-reasoning", + "reply-mixed-thinking-2-reasoning", + "reply-mixed-thinking-2-reasoning"), + messageIds); + } + @Test void testThinkingEventsAreIgnoredWhenReasoningDisabled() { List events = From 7bdcea2fc73b11b49bcdbb0d997282649c3394c2 Mon Sep 17 00:00:00 2001 From: jujn <2087687391@qq.com> Date: Tue, 8 Sep 2026 19:44:56 +0800 Subject: [PATCH 3/3] fix: improve --- .../adapter/strategy/AguiStreamContext.java | 97 ++----------------- .../strategy/TextBlockEventConverter.java | 11 ++- .../strategy/ThinkingBlockEventConverter.java | 12 ++- .../agui/adapter/AguiAgentAdapterV2Test.java | 41 ++++---- 4 files changed, 43 insertions(+), 118 deletions(-) diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java index 039e7b0eec..086db9b6b0 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/AguiStreamContext.java @@ -38,10 +38,6 @@ public class AguiStreamContext { - // CopilotKit will merge reasoning and text with the same messageId, adding suffixes to the - // reasoning to distinguish them - public static final String REASONING_MESSAGE_ID_SUFFIX = "-reasoning"; - private static final Logger logger = LoggerFactory.getLogger(AguiStreamContext.class); private final String threadId; @@ -53,12 +49,8 @@ public class AguiStreamContext { private final Set startedTextMessages = new LinkedHashSet<>(); private final Set endedTextMessages = new LinkedHashSet<>(); - private final Map textMessageIds = new LinkedHashMap<>(); - private final Map firstTextBlockByReply = new LinkedHashMap<>(); private final Set startedReasoningMessages = new LinkedHashSet<>(); private final Set endedReasoningMessages = new LinkedHashSet<>(); - private final Map reasoningMessageIds = new LinkedHashMap<>(); - private final Map firstThinkingBlockByReply = new LinkedHashMap<>(); private final Set startedToolCalls = new LinkedHashSet<>(); private final Set endedToolCalls = new LinkedHashSet<>(); private String currentTextMessageId; @@ -152,26 +144,6 @@ public void appendTextDelta(String messageId, String delta) { } } - public String textMessageId(String replyId, String blockId) { - String normalizedBlockId = isBlank(blockId) ? "text" : blockId; - TextBlockKey key = new TextBlockKey(replyId, normalizedBlockId); - return textMessageIds.computeIfAbsent( - key, - ignored -> { - String firstBlockId = - firstTextBlockByReply.putIfAbsent(replyId, normalizedBlockId); - if (firstBlockId == null || Objects.equals(firstBlockId, normalizedBlockId)) { - return replyId; - } - return replyId + "-" + normalizedBlockId; - }); - } - - public String existingTextMessageId(String replyId, String blockId) { - String normalizedBlockId = isBlank(blockId) ? "text" : blockId; - return textMessageIds.get(new TextBlockKey(replyId, normalizedBlockId)); - } - public void closeActiveTextMessage() { if (currentTextMessageId == null) { return; @@ -193,50 +165,19 @@ public void closeTextMessage(String messageId) { } public void startReasoningMessage(String messageId) { - String reasoningMessageId = withReasoningSuffix(messageId); - if (startedReasoningMessages.add(reasoningMessageId)) { - emit( - new AguiEvent.ReasoningMessageStart( - threadId, runId, reasoningMessageId, "reasoning")); + if (startedReasoningMessages.add(messageId)) { + emit(new AguiEvent.ReasoningMessageStart(threadId, runId, messageId, "reasoning")); } - currentReasoningMessageId = reasoningMessageId; + currentReasoningMessageId = messageId; } public void appendReasoningDelta(String messageId, String delta) { if (delta != null && !delta.isEmpty()) { startReasoningMessage(messageId); - emit( - new AguiEvent.ReasoningMessageContent( - threadId, runId, withReasoningSuffix(messageId), delta)); + emit(new AguiEvent.ReasoningMessageContent(threadId, runId, messageId, delta)); } } - public String reasoningMessageId(String replyId, String blockId) { - String normalizedBlockId = isBlank(blockId) ? "thinking" : blockId; - ThinkingBlockKey key = new ThinkingBlockKey(replyId, normalizedBlockId); - String messageId = - reasoningMessageIds.computeIfAbsent( - key, - ignored -> { - String firstBlockId = - firstThinkingBlockByReply.putIfAbsent( - replyId, normalizedBlockId); - if (firstBlockId == null - || Objects.equals(firstBlockId, normalizedBlockId)) { - return replyId; - } - return replyId + "-" + normalizedBlockId; - }); - return withReasoningSuffix(messageId); - } - - public String existingReasoningMessageId(String replyId, String blockId) { - String normalizedBlockId = isBlank(blockId) ? "thinking" : blockId; - String messageId = - reasoningMessageIds.get(new ThinkingBlockKey(replyId, normalizedBlockId)); - return messageId == null ? null : withReasoningSuffix(messageId); - } - public void closeActiveReasoningMessage() { if (currentReasoningMessageId == null) { return; @@ -245,17 +186,16 @@ public void closeActiveReasoningMessage() { } public void closeReasoningMessage(String messageId) { - String reasoningMessageId = withReasoningSuffix(messageId); - if (reasoningMessageId == null - || !startedReasoningMessages.contains(reasoningMessageId) - || endedReasoningMessages.contains(reasoningMessageId)) { + if (messageId == null + || !startedReasoningMessages.contains(messageId) + || endedReasoningMessages.contains(messageId)) { return; } - endedReasoningMessages.add(reasoningMessageId); - if (Objects.equals(reasoningMessageId, currentReasoningMessageId)) { + endedReasoningMessages.add(messageId); + if (Objects.equals(messageId, currentReasoningMessageId)) { currentReasoningMessageId = null; } - emit(new AguiEvent.ReasoningMessageEnd(threadId, runId, reasoningMessageId)); + emit(new AguiEvent.ReasoningMessageEnd(threadId, runId, messageId)); } public void startToolCall(String toolCallId, String toolCallName) { @@ -329,7 +269,6 @@ public void endToolResult(String replyId, String toolCallId) { if (endedToolCalls.add(toolCallId)) { emit(new AguiEvent.ToolCallEnd(threadId, runId, toolCallId)); } - StringBuilder content = toolResultContent.remove(toolCallId); emit( new AguiEvent.ToolCallResult( @@ -385,16 +324,6 @@ private static String normalizeToolCallName(String toolCallName) { return toolCallName != null && !toolCallName.isBlank() ? toolCallName : "unknown"; } - private static String withReasoningSuffix(String messageId) { - if (messageId == null) { - return null; - } - if (messageId.endsWith(REASONING_MESSAGE_ID_SUFFIX)) { - return messageId; - } - return messageId + REASONING_MESSAGE_ID_SUFFIX; - } - private static String serialize(ContentBlock data) { if (data instanceof TextBlock textBlock) { return textBlock.getText(); @@ -410,10 +339,6 @@ private static boolean isBlank(String value) { return value == null || value.isBlank(); } - private record TextBlockKey(String replyId, String blockId) {} - - private record ThinkingBlockKey(String replyId, String blockId) {} - private void warnMissingToolCallId(String eventName) { if (!warnedMissingToolCallIdOperations.add(eventName)) { return; @@ -455,7 +380,6 @@ private static Set frontendToolNames(RunAgentInput runInput) { } static final class TokenUsageAccumulator { - private long cumulativeInputTokens; private long cumulativeOutputTokens; private long cumulativeCachedTokens; @@ -483,7 +407,6 @@ TokenUsageSnapshot add(ChatUsage usage) { record TokenUsageSnapshot(TokenUsage delta, TokenUsage cumulative) {} record TokenUsage(long inputTokens, long outputTokens, long cachedTokens, double time) { - long totalTokens() { return inputTokens + outputTokens; } diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/TextBlockEventConverter.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/TextBlockEventConverter.java index 6980b4343f..49e1d5c8d9 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/TextBlockEventConverter.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/TextBlockEventConverter.java @@ -33,11 +33,14 @@ public Set> eventTypes() { public void convert(AgentEvent event, AguiStreamContext context) { if (event instanceof TextBlockDeltaEvent delta) { // AguiEvent.TextMessageStart delays sending when content arrives - String messageId = context.textMessageId(delta.getReplyId(), delta.getBlockId()); - context.appendTextDelta(messageId, delta.getDelta()); + context.appendTextDelta( + messageId(delta.getReplyId(), delta.getBlockId()), delta.getDelta()); } else if (event instanceof TextBlockEndEvent end) { - String messageId = context.existingTextMessageId(end.getReplyId(), end.getBlockId()); - context.closeTextMessage(messageId); + context.closeTextMessage(messageId(end.getReplyId(), end.getBlockId())); } } + + private String messageId(String replyId, String blockId) { + return replyId + "-" + blockId; + } } diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ThinkingBlockEventConverter.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ThinkingBlockEventConverter.java index 12e08b2bb6..0994d2e041 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ThinkingBlockEventConverter.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ThinkingBlockEventConverter.java @@ -40,12 +40,14 @@ public void convert(AgentEvent event, AguiStreamContext context) { if (event instanceof ThinkingBlockDeltaEvent delta) { // AguiEvent.ReasoningMessageStart delays sending when content arrives - String messageId = context.reasoningMessageId(delta.getReplyId(), delta.getBlockId()); - context.appendReasoningDelta(messageId, delta.getDelta()); + context.appendReasoningDelta( + messageId(delta.getReplyId(), delta.getBlockId()), delta.getDelta()); } else if (event instanceof ThinkingBlockEndEvent end) { - String messageId = - context.existingReasoningMessageId(end.getReplyId(), end.getBlockId()); - context.closeReasoningMessage(messageId); + context.closeReasoningMessage(messageId(end.getReplyId(), end.getBlockId())); } } + + private String messageId(String replyId, String blockId) { + return replyId + "-" + blockId; + } } diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java index c8bf95b823..421a390b9f 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterV2Test.java @@ -374,9 +374,9 @@ void testTextSegmentsSeparatedByToolCallUseDistinctMessageIds() { .toList(); assertEquals( List.of( - "reply-mixed", - "reply-mixed", - "reply-mixed", + "reply-mixed-text", + "reply-mixed-text", + "reply-mixed-text", "reply-mixed-text-2", "reply-mixed-text-2", "reply-mixed-text-2"), @@ -436,12 +436,12 @@ void testReasoningSegmentsSeparatedByToolCallUseDistinctMessageIds() { .toList(); assertEquals( List.of( - "reply-mixed-reasoning", - "reply-mixed-reasoning", - "reply-mixed-reasoning", - "reply-mixed-thinking-2-reasoning", - "reply-mixed-thinking-2-reasoning", - "reply-mixed-thinking-2-reasoning"), + "reply-mixed-thinking", + "reply-mixed-thinking", + "reply-mixed-thinking", + "reply-mixed-thinking-2", + "reply-mixed-thinking-2", + "reply-mixed-thinking-2"), messageIds); } @@ -463,9 +463,9 @@ void testThinkingEventsConvertWhenReasoningEnabled() { List events = runReActEvents( AguiAdapterConfig.builder().enableReasoning(true).build(), - new ThinkingBlockStartEvent("reply-thinking", "block-1"), - new ThinkingBlockDeltaEvent("reply-thinking", "block-1", "visible"), - new ThinkingBlockEndEvent("reply-thinking", "block-1")); + new ThinkingBlockStartEvent("reply-thinking", "thinking"), + new ThinkingBlockDeltaEvent("reply-thinking", "thinking", "visible"), + new ThinkingBlockEndEvent("reply-thinking", "thinking")); assertEquals( List.of( @@ -479,8 +479,7 @@ void testThinkingEventsConvertWhenReasoningEnabled() { assertInstanceOf(AguiEvent.ReasoningMessageContent.class, events.get(1)); AguiEvent.ReasoningMessageEnd end = assertInstanceOf(AguiEvent.ReasoningMessageEnd.class, events.get(2)); - String expectedMessageId = - "reply-thinking" + AguiStreamContext.REASONING_MESSAGE_ID_SUFFIX; + String expectedMessageId = "reply-thinking-thinking"; assertEquals(expectedMessageId, start.messageId()); assertEquals(expectedMessageId, content.messageId()); assertEquals(expectedMessageId, end.messageId()); @@ -491,10 +490,10 @@ void testTextAndReasoningUseDifferentMessageIdsForSameReply() { List events = runReActEvents( AguiAdapterConfig.builder().enableReasoning(true).build(), - new ThinkingBlockDeltaEvent("reply-shared", "thinking-1", "think"), - new ThinkingBlockEndEvent("reply-shared", "thinking-1"), - new TextBlockDeltaEvent("reply-shared", "text-1", "answer"), - new TextBlockEndEvent("reply-shared", "text-1")); + new ThinkingBlockDeltaEvent("reply-shared", "thinking", "think"), + new ThinkingBlockEndEvent("reply-shared", "thinking"), + new TextBlockDeltaEvent("reply-shared", "text", "answer"), + new TextBlockEndEvent("reply-shared", "text")); AguiEvent.ReasoningMessageContent reasoningContent = events.stream() @@ -509,10 +508,8 @@ void testTextAndReasoningUseDifferentMessageIdsForSameReply() { .findFirst() .orElseThrow(); - assertEquals("reply-shared", textContent.messageId()); - assertEquals( - "reply-shared" + AguiStreamContext.REASONING_MESSAGE_ID_SUFFIX, - reasoningContent.messageId()); + assertEquals("reply-shared-text", textContent.messageId()); + assertEquals("reply-shared-thinking", reasoningContent.messageId()); } @Test