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 75f753e4c9..34f5bef62b 100644 --- a/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java +++ b/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java @@ -51,6 +51,7 @@ import io.agentscope.core.event.ToolCallDeltaEvent; import io.agentscope.core.event.ToolCallEndEvent; import io.agentscope.core.event.ToolCallStartEvent; +import io.agentscope.core.event.ToolProgressEvent; import io.agentscope.core.event.ToolResultDataDeltaEvent; import io.agentscope.core.event.ToolResultEndEvent; import io.agentscope.core.event.ToolResultStartEvent; @@ -3736,47 +3737,25 @@ private Flux runToolBatch( tool.getName())); } - Set chunkedToolIds = - ConcurrentHashMap.newKeySet(); - BiConsumer internalChunkCallback = (toolUse, chunk) -> { if (chunk.getOutput() != null && !chunk.getOutput() .isEmpty()) { - chunkedToolIds.add( - toolUse.getId()); for (ContentBlock block : chunk.getOutput()) { - if (block - instanceof - TextBlock tb) { - sink.next( - new ToolResultTextDeltaEvent( - replyId, - toolUse - .getId(), - toolUse - .getName(), - tb - .getText()) - .withMetadata( - chunk - .getMetadata())); - } else { - sink.next( - new ToolResultDataDeltaEvent( - replyId, - toolUse - .getId(), - toolUse - .getName(), - block) - .withMetadata( - chunk - .getMetadata())); - } + sink.next( + new ToolProgressEvent( + replyId, + toolUse + .getId(), + toolUse + .getName(), + block) + .withMetadata( + chunk + .getMetadata())); } } hookDispatcher @@ -3854,8 +3833,7 @@ private Flux runToolBatch( emitToolResultDelta( sink, replyId, - entry, - chunkedToolIds); + entry); ToolResultState state = determineToolResultState( @@ -3868,7 +3846,10 @@ private Flux runToolBatch( .getId(), entry.getKey() .getName(), - state) + state, + finalToolResultText( + entry + .getValue())) .withMetadata( entry.getValue() .getMetadata())); @@ -3993,21 +3974,20 @@ private List getSuspendedToolCalls( } /** - * Emit delta events for tool results that were NOT already streamed via the chunk - * callback. For non-streaming tools the chunk callback is never invoked, so the - * event stream would otherwise contain only START and END with no content. + * Emit delta events carrying the tool method's return value. + * + *

Progress published through {@link io.agentscope.core.tool.ToolEmitter} travels as + * {@link ToolProgressEvent} rather than as deltas, so every delta here is part of what the + * model received. A tool that streams progress and also returns a value therefore reports + * both: its progress chunks on their own type, then its result chunks. */ private void emitToolResultDelta( FluxSink sink, String replyId, - Map.Entry entry, - Set chunkedToolIds) { + Map.Entry entry) { String toolId = entry.getKey().getId(); String toolName = entry.getKey().getName(); ToolResultBlock toolResult = entry.getValue(); - if (chunkedToolIds.contains(toolId)) { - return; - } List output = toolResult.getOutput(); if (output == null || output.isEmpty()) { return; @@ -4025,6 +4005,34 @@ private void emitToolResultDelta( } } + /** + * Join the text blocks of a tool method's return value, reported on {@link + * ToolResultEndEvent#getFinalResultText()}. + * + *

The deltas of this lifecycle already carry the same value, and progress published + * through {@link io.agentscope.core.tool.ToolEmitter} no longer reaches them; this field is + * what a consumer that does not accumulate the stream still gets the result from, in a + * single event. + * + *

An empty join is reported as {@code ""}, not {@code null}: a tool that returns an image + * or a blank string did produce a return value, and a consumer reading {@code null} as + * "nothing was reported" would go on looking for it elsewhere. Only a result we know nothing + * about is left unreported. + * + * @return joined text (possibly empty), or {@code null} when there are no content blocks to + * report from + */ + private String finalToolResultText(ToolResultBlock result) { + if (result == null || result.getOutput() == null || result.getOutput().isEmpty()) { + return null; + } + return result.getOutput().stream() + .filter(TextBlock.class::isInstance) + .map(block -> ((TextBlock) block).getText()) + .filter(value -> value != null && !value.isEmpty()) + .collect(Collectors.joining()); + } + private ToolResultState determineToolResultState(ToolResultBlock result) { if (result.isSuspended()) { return ToolResultState.SUSPENDED; diff --git a/agentscope-core/src/main/java/io/agentscope/core/event/AgentEvent.java b/agentscope-core/src/main/java/io/agentscope/core/event/AgentEvent.java index 60b9b8b4ec..711e195002 100644 --- a/agentscope-core/src/main/java/io/agentscope/core/event/AgentEvent.java +++ b/agentscope-core/src/main/java/io/agentscope/core/event/AgentEvent.java @@ -54,6 +54,7 @@ @JsonSubTypes.Type(value = ToolCallDeltaEvent.class, name = "TOOL_CALL_DELTA"), @JsonSubTypes.Type(value = ToolCallEndEvent.class, name = "TOOL_CALL_END"), @JsonSubTypes.Type(value = ToolResultStartEvent.class, name = "TOOL_RESULT_START"), + @JsonSubTypes.Type(value = ToolProgressEvent.class, name = "TOOL_PROGRESS"), @JsonSubTypes.Type(value = ToolResultTextDeltaEvent.class, name = "TOOL_RESULT_TEXT_DELTA"), @JsonSubTypes.Type(value = ToolResultDataDeltaEvent.class, name = "TOOL_RESULT_DATA_DELTA"), @JsonSubTypes.Type(value = ToolResultEndEvent.class, name = "TOOL_RESULT_END"), diff --git a/agentscope-core/src/main/java/io/agentscope/core/event/AgentEventType.java b/agentscope-core/src/main/java/io/agentscope/core/event/AgentEventType.java index 69cd2f046f..bf3cfd1d1b 100644 --- a/agentscope-core/src/main/java/io/agentscope/core/event/AgentEventType.java +++ b/agentscope-core/src/main/java/io/agentscope/core/event/AgentEventType.java @@ -69,6 +69,7 @@ public enum AgentEventType { TOOL_CALL_END("TOOL_CALL_END"), TOOL_RESULT_START("TOOL_RESULT_START"), + TOOL_PROGRESS("TOOL_PROGRESS"), TOOL_RESULT_TEXT_DELTA("TOOL_RESULT_TEXT_DELTA"), @JsonAlias({"TOOL_RESULT_BINARY_DELTA"}) TOOL_RESULT_DATA_DELTA("TOOL_RESULT_DATA_DELTA"), diff --git a/agentscope-core/src/main/java/io/agentscope/core/event/ToolProgressEvent.java b/agentscope-core/src/main/java/io/agentscope/core/event/ToolProgressEvent.java new file mode 100644 index 0000000000..2de228032e --- /dev/null +++ b/agentscope-core/src/main/java/io/agentscope/core/event/ToolProgressEvent.java @@ -0,0 +1,116 @@ +/* + * 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.core.event; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; +import io.agentscope.core.message.ContentBlock; +import io.agentscope.core.message.TextBlock; +import java.util.Map; + +/** + * An intermediate update published by a running tool through {@link + * io.agentscope.core.tool.ToolEmitter}. + * + *

{@link ToolResultTextDeltaEvent} and {@link ToolResultDataDeltaEvent} carry the tool's result, + * which is also what the model receives. Progress carries no such promise: {@code ToolEmitter}'s + * contract says emitted chunks are not sent to the LLM, so a consumer that rebuilt a tool result + * from the delta stream would persist progress as the return value. Keeping the two on separate + * event types makes that confusion impossible rather than merely discouraged. + */ +public class ToolProgressEvent extends AgentEvent { + + private final String replyId; + private final String toolCallId; + private final String toolCallName; + private final ContentBlock content; + + @JsonCreator + public ToolProgressEvent( + @JsonProperty("id") String id, + @JsonProperty("createdAt") String createdAt, + @JsonProperty("replyId") String replyId, + @JsonProperty("toolCallId") String toolCallId, + @JsonProperty("toolCallName") String toolCallName, + @JsonProperty("content") ContentBlock content, + @JsonProperty("metadata") Map metadata) { + super(id, createdAt); + this.replyId = replyId; + this.toolCallId = toolCallId; + this.toolCallName = toolCallName; + this.content = content; + this.withMetadata(metadata); + } + + /** + * Backward-compatible constructor for callers that do not provide metadata. + */ + public ToolProgressEvent( + String id, + String createdAt, + String replyId, + String toolCallId, + String toolCallName, + ContentBlock content) { + this(id, createdAt, replyId, toolCallId, toolCallName, content, null); + } + + public ToolProgressEvent( + String replyId, String toolCallId, String toolCallName, ContentBlock content) { + this.replyId = replyId; + this.toolCallId = toolCallId; + this.toolCallName = toolCallName; + this.content = content; + } + + @Override + public AgentEventType getType() { + return AgentEventType.TOOL_PROGRESS; + } + + public String getReplyId() { + return replyId; + } + + public String getToolCallId() { + return toolCallId; + } + + public String getToolCallName() { + return toolCallName; + } + + /** + * The progress chunk as published by the tool. + * + * @return the emitted content block, text or otherwise + */ + public ContentBlock getContent() { + return content; + } + + /** + * The progress chunk when it is text. + * + *

Most consumers only render progress, so the common case is exposed directly; a non-text + * chunk reads as {@code null} here and stays available through {@link #getContent()}. + * + * @return the chunk's text, or {@code null} when the chunk is not a text block + */ + public String getText() { + return content instanceof TextBlock textBlock ? textBlock.getText() : null; + } +} diff --git a/agentscope-core/src/main/java/io/agentscope/core/event/ToolResultEndEvent.java b/agentscope-core/src/main/java/io/agentscope/core/event/ToolResultEndEvent.java index b3ff00dd50..9da90992e5 100644 --- a/agentscope-core/src/main/java/io/agentscope/core/event/ToolResultEndEvent.java +++ b/agentscope-core/src/main/java/io/agentscope/core/event/ToolResultEndEvent.java @@ -27,6 +27,19 @@ public class ToolResultEndEvent extends AgentEvent { private final String toolCallName; private final ToolResultState state; + /** + * The tool method's return value as text, or {@code null} when the producer reports nothing. + * + *

{@code ""} is a reported value, not an absent one: it means the tool returned content with + * no text in it (an image, or a blank result), and a consumer must not fall back to its delta + * buffer for that case. + * + *

This deliberately does not ride on {@link AgentEvent#getMetadata()}, which passes the tool + * result's own metadata through unchanged — {@code ReActAgentNewLoopE2ETest} asserts that map + * exactly, so a framework key in it would be a contract change for unrelated consumers. + */ + private final String finalResultText; + @JsonCreator public ToolResultEndEvent( @JsonProperty("id") String id, @@ -35,12 +48,14 @@ public ToolResultEndEvent( @JsonProperty("toolCallId") String toolCallId, @JsonProperty("toolCallName") String toolCallName, @JsonProperty("state") ToolResultState state, - @JsonProperty("metadata") Map metadata) { + @JsonProperty("metadata") Map metadata, + @JsonProperty("finalResultText") String finalResultText) { super(id, createdAt); this.replyId = replyId; this.toolCallId = toolCallId; this.toolCallName = toolCallName; this.state = state; + this.finalResultText = finalResultText; this.withMetadata(metadata); } @@ -54,15 +69,29 @@ public ToolResultEndEvent( String toolCallId, String toolCallName, ToolResultState state) { - this(id, createdAt, replyId, toolCallId, toolCallName, state, null); + this(id, createdAt, replyId, toolCallId, toolCallName, state, null, null); } public ToolResultEndEvent( String replyId, String toolCallId, String toolCallName, ToolResultState state) { + this(replyId, toolCallId, toolCallName, state, null); + } + + /** + * As {@link #ToolResultEndEvent(String, String, String, ToolResultState)} additionally reporting + * the tool method's return value. + */ + public ToolResultEndEvent( + String replyId, + String toolCallId, + String toolCallName, + ToolResultState state, + String finalResultText) { this.replyId = replyId; this.toolCallId = toolCallId; this.toolCallName = toolCallName; this.state = state; + this.finalResultText = finalResultText; } @Override @@ -85,4 +114,14 @@ public String getToolCallName() { public ToolResultState getState() { return state; } + + /** + * The tool method's return value as text. + * + * @return the return value (possibly empty), or {@code null} when the producer reports nothing, + * which leaves consumers falling back to the delta stream + */ + public String getFinalResultText() { + return finalResultText; + } } diff --git a/agentscope-core/src/main/java/io/agentscope/core/tool/ToolEmitter.java b/agentscope-core/src/main/java/io/agentscope/core/tool/ToolEmitter.java index 0a299420c3..6dbd76079e 100644 --- a/agentscope-core/src/main/java/io/agentscope/core/tool/ToolEmitter.java +++ b/agentscope-core/src/main/java/io/agentscope/core/tool/ToolEmitter.java @@ -22,8 +22,11 @@ * *

Tool methods can declare a ToolEmitter parameter to send intermediate progress updates and * messages during execution. These streaming chunks are delivered to registered hooks via - * {@code onActingChunk()} events but are NOT sent to the LLM. Only the final return value of the - * tool method is sent to the LLM as the tool result. + * {@code onActingChunk()} events and to event-stream subscribers as {@link + * io.agentscope.core.event.ToolProgressEvent}, but are NOT sent to the LLM. Only the final return + * value of the tool method is sent to the LLM as the tool result, which is what {@link + * io.agentscope.core.event.ToolResultTextDeltaEvent} and {@link + * io.agentscope.core.event.ToolResultDataDeltaEvent} carry. * *

Key Characteristics: *

    diff --git a/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopE2ETest.java b/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopE2ETest.java index 2220f562b7..843647f139 100644 --- a/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopE2ETest.java +++ b/agentscope-core/src/test/java/io/agentscope/core/agent/ReActAgentNewLoopE2ETest.java @@ -26,15 +26,18 @@ import io.agentscope.core.event.AgentStartEvent; import io.agentscope.core.event.ModelCallEndEvent; import io.agentscope.core.event.ToolCallEndEvent; +import io.agentscope.core.event.ToolProgressEvent; import io.agentscope.core.event.ToolResultEndEvent; import io.agentscope.core.event.ToolResultTextDeltaEvent; import io.agentscope.core.message.ContentBlock; +import io.agentscope.core.message.ImageBlock; import io.agentscope.core.message.Msg; import io.agentscope.core.message.MsgRole; import io.agentscope.core.message.TextBlock; import io.agentscope.core.message.ToolResultBlock; import io.agentscope.core.message.ToolResultState; import io.agentscope.core.message.ToolUseBlock; +import io.agentscope.core.message.URLSource; import io.agentscope.core.middleware.ActingInput; import io.agentscope.core.middleware.AgentInput; import io.agentscope.core.middleware.MiddlewareBase; @@ -416,4 +419,116 @@ void streamEventsRestoresEmitterWhenActingContextLosesEventKeys() { assertTrue(events.get(0) instanceof AgentStartEvent); assertTrue(events.get(events.size() - 1) instanceof AgentEndEvent); } + + @Test + void progressAndResultTravelOnTheirOwnEventTypes() { + List events = runWithProgressTool(ToolResultBlock.text("fingerprint enrolled")); + + List progress = + events.stream() + .filter(ToolProgressEvent.class::isInstance) + .map(ToolProgressEvent.class::cast) + .map(ToolProgressEvent::getText) + .toList(); + assertTrue( + progress.contains("scanning 3 dirs"), + "the emitted chunk has to reach the stream as progress"); + + List resultDeltas = + events.stream() + .filter(ToolResultTextDeltaEvent.class::isInstance) + .map(e -> ((ToolResultTextDeltaEvent) e).getDelta()) + .toList(); + assertEquals( + List.of("fingerprint enrolled"), + resultDeltas, + "deltas carry only what the model received, so progress must not appear here"); + + assertEquals( + "fingerprint enrolled", + lastToolResultEnd(events).getFinalResultText(), + "the end event reports the same return value in one event"); + } + + @Test + void imageOnlyReturnValueIsReportedAsEmptyRatherThanAbsent() { + // "" and null mean different things downstream: null reads as "the producer reported + // nothing", so a consumer would keep looking for a value the tool did return. A tool that + // returned only an image produced a value, and that value has no text in it. + List events = + runWithProgressTool( + ToolResultBlock.of( + ImageBlock.builder() + .source(new URLSource("https://example.com/chart.png")) + .build())); + + assertEquals("", lastToolResultEnd(events).getFinalResultText()); + } + + private List runWithProgressTool(ToolResultBlock returnValue) { + ScriptedModel model = + new ScriptedModel( + List.of( + () -> Flux.just(toolUseResponse("c1", "enroll", "alpha")), + () -> Flux.just(textResponse("done")))); + Toolkit tk = new Toolkit(); + tk.registerAgentTool(new ProgressTool("enroll", returnValue)); + + List events = + ReActAgent.builder() + .name("asst") + .sysPrompt("you are helpful") + .model(model) + .toolkit(tk) + .build() + .streamEvents( + List.of( + Msg.builder() + .role(MsgRole.USER) + .textContent("run enroll") + .build())) + .collectList() + .block(); + assertNotNull(events); + return events; + } + + private static ToolResultEndEvent lastToolResultEnd(List events) { + return events.stream() + .filter(ToolResultEndEvent.class::isInstance) + .map(ToolResultEndEvent.class::cast) + .reduce((first, second) -> second) + .orElseThrow(); + } + + /** Streams one progress chunk, then returns {@code returnValue}: the {@code ToolEmitter} deal. */ + private static final class ProgressTool extends ToolBase { + private final ToolResultBlock returnValue; + + ProgressTool(String name, ToolResultBlock returnValue) { + super( + name, + "streams progress then returns", + AlwaysAllowTool.schema(), + true, + true, + false, + null, + false, + false); + this.returnValue = returnValue; + } + + @Override + public Mono checkPermissions( + Map input, PermissionContextState ctx) { + return Mono.just(PermissionDecision.allow("ok")); + } + + @Override + public Mono callAsync(ToolCallParam param) { + param.getEmitter().emit(ToolResultBlock.text("scanning 3 dirs")); + return Mono.just(returnValue); + } + } } diff --git a/agentscope-core/src/test/java/io/agentscope/core/event/AgentEventStreamTest.java b/agentscope-core/src/test/java/io/agentscope/core/event/AgentEventStreamTest.java index 5adf5976fe..46caf884ae 100644 --- a/agentscope-core/src/test/java/io/agentscope/core/event/AgentEventStreamTest.java +++ b/agentscope-core/src/test/java/io/agentscope/core/event/AgentEventStreamTest.java @@ -16,7 +16,9 @@ package io.agentscope.core.event; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -161,6 +163,31 @@ void toolResultTextDeltaEventMetadataRoundTrip() throws Exception { assertEquals(metadata, deserialized.getMetadata()); } + @Test + @DisplayName("ToolProgressEvent round-trips its type, content and metadata") + void toolProgressEventRoundTrip() throws Exception { + Map metadata = Map.of("chunk", true); + ToolProgressEvent original = + new ToolProgressEvent( + "reply-1", + "tc-1", + "search", + TextBlock.builder().text("scanning 3 dirs").build()); + original.withMetadata(metadata); + + String json = mapper.writeValueAsString(original); + assertTrue(json.contains("TOOL_PROGRESS"), json); + + AgentEvent deserialized = mapper.readValue(json, AgentEvent.class); + assertTrue( + deserialized instanceof ToolProgressEvent, + "the subtype registry has to resolve progress events"); + ToolProgressEvent back = (ToolProgressEvent) deserialized; + assertEquals("tc-1", back.getToolCallId()); + assertEquals("scanning 3 dirs", back.getText()); + assertEquals(metadata, back.getMetadata()); + } + @Test @DisplayName("ToolResultDataDeltaEvent serializes and deserializes metadata") void toolResultDataDeltaEventMetadataRoundTrip() throws Exception { @@ -190,6 +217,61 @@ void toolResultStartEventEmptyMetadata() throws Exception { deserialized.getMetadata() == null || deserialized.getMetadata().isEmpty(), "ToolResultStartEvent should not carry result metadata before execution"); } + + @Test + @DisplayName("finalResultText survives the polymorphic round trip") + void finalResultTextSurvivesRoundTrip() throws Exception { + ToolResultEndEvent original = + new ToolResultEndEvent( + "reply-1", + "tc-1", + "search", + ToolResultState.SUCCESS, + "fingerprint enrolled"); + + String json = mapper.writeValueAsString(original); + assertTrue( + json.contains("\"finalResultText\":\"fingerprint enrolled\""), + "the new property has to be on the wire, not just in the object: " + json); + + ToolResultEndEvent back = + assertInstanceOf( + ToolResultEndEvent.class, mapper.readValue(json, AgentEvent.class)); + assertEquals("fingerprint enrolled", back.getFinalResultText()); + } + + @Test + @DisplayName("an empty return value decodes as empty, not as absent") + void emptyFinalResultTextIsNotLossy() throws Exception { + // "" carries meaning: the tool did return content, it just holds no text. Read back as + // null it would send consumers falling back to the delta buffer, re-opening the leak. + ToolResultEndEvent original = + new ToolResultEndEvent( + "reply-1", "tc-1", "search", ToolResultState.SUCCESS, ""); + + ToolResultEndEvent back = + assertInstanceOf( + ToolResultEndEvent.class, + mapper.readValue( + mapper.writeValueAsString(original), AgentEvent.class)); + assertEquals("", back.getFinalResultText()); + } + + @Test + @DisplayName("a payload written before the field existed decodes to null") + void absentFinalResultTextDecodesToNull() throws Exception { + // Sessions persisted by an older build are reloaded through this same creator. + ToolResultEndEvent back = + assertInstanceOf( + ToolResultEndEvent.class, + mapper.readValue( + "{\"type\":\"TOOL_RESULT_END\",\"replyId\":\"r\"," + + "\"toolCallId\":\"t\",\"toolCallName\":\"n\"," + + "\"state\":\"success\"}", + AgentEvent.class)); + assertEquals("r", back.getReplyId()); + assertNull(back.getFinalResultText()); + } } @Nested diff --git a/agentscope-core/src/test/java/io/agentscope/core/tool/ReActAgentToolkitConcurrencyTest.java b/agentscope-core/src/test/java/io/agentscope/core/tool/ReActAgentToolkitConcurrencyTest.java index d83c0b902b..50a5ac7727 100644 --- a/agentscope-core/src/test/java/io/agentscope/core/tool/ReActAgentToolkitConcurrencyTest.java +++ b/agentscope-core/src/test/java/io/agentscope/core/tool/ReActAgentToolkitConcurrencyTest.java @@ -23,7 +23,7 @@ import io.agentscope.core.ReActAgent; import io.agentscope.core.agent.RuntimeContext; import io.agentscope.core.event.AgentEvent; -import io.agentscope.core.event.ToolResultTextDeltaEvent; +import io.agentscope.core.event.ToolProgressEvent; import io.agentscope.core.message.Msg; import io.agentscope.core.message.MsgRole; import io.agentscope.core.message.TextBlock; @@ -323,12 +323,12 @@ void concurrentCallsDoNotCrossTalkChunkCallbacks() { Flux streamB = agent.streamEvents(List.of(userMsg("prompt-B")), rcB); Mono> deltasA = - streamA.ofType(ToolResultTextDeltaEvent.class) - .map(ToolResultTextDeltaEvent::getDelta) + streamA.ofType(ToolProgressEvent.class) + .map(ToolProgressEvent::getText) .collectList(); Mono> deltasB = - streamB.ofType(ToolResultTextDeltaEvent.class) - .map(ToolResultTextDeltaEvent::getDelta) + streamB.ofType(ToolProgressEvent.class) + .map(ToolProgressEvent::getText) .collectList(); var tuple = Mono.zip(deltasA, deltasB).block(TIMEOUT); diff --git a/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/hitl/InterruptionExample.java b/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/hitl/InterruptionExample.java index fa74fa34e7..6e3c594bdf 100644 --- a/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/hitl/InterruptionExample.java +++ b/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/hitl/InterruptionExample.java @@ -19,7 +19,7 @@ import io.agentscope.core.agent.Agent; import io.agentscope.core.agent.RuntimeContext; import io.agentscope.core.event.AgentEvent; -import io.agentscope.core.event.ToolResultTextDeltaEvent; +import io.agentscope.core.event.ToolProgressEvent; import io.agentscope.core.message.Msg; import io.agentscope.core.message.ToolResultBlock; import io.agentscope.core.message.UserMessage; @@ -46,7 +46,7 @@ * {@code next.apply()}. *
  • {@code PreActingEvent} / {@code PostActingEvent} → {@code onActing()} before/after * {@code next.apply()}.
  • - *
  • {@code ActingChunkEvent} → tap {@code ToolResultTextDeltaEvent} in acting stream.
  • + *
  • {@code ActingChunkEvent} → tap {@code ToolProgressEvent} in acting stream.
  • *
  • {@code ErrorEvent} → {@code doOnError()} on the agent stream.
  • *
  • {@code agent.getMemory().getMessages()} → {@code agent.getAgentState().getContext()}.
  • *
  • Removed {@code .memory(new InMemoryMemory())}.
  • @@ -189,9 +189,9 @@ public Flux onActing( return next.apply(input) .doOnNext( event -> { - if (event instanceof ToolResultTextDeltaEvent delta) { + if (event instanceof ToolProgressEvent progress) { System.out.println( - "[Middleware] Tool progress: " + delta.getDelta()); + "[Middleware] Tool progress: " + progress.getText()); } }) .doOnComplete(() -> System.out.println("[Middleware] Tool result completed")); diff --git a/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/middleware/CustomizedMiddlewareExample.java b/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/middleware/CustomizedMiddlewareExample.java index d026bed0d7..15abe46e98 100644 --- a/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/middleware/CustomizedMiddlewareExample.java +++ b/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/middleware/CustomizedMiddlewareExample.java @@ -20,7 +20,7 @@ import io.agentscope.core.agent.RuntimeContext; import io.agentscope.core.event.AgentEvent; import io.agentscope.core.event.TextBlockDeltaEvent; -import io.agentscope.core.event.ToolResultTextDeltaEvent; +import io.agentscope.core.event.ToolProgressEvent; import io.agentscope.core.message.Msg; import io.agentscope.core.message.ToolResultBlock; import io.agentscope.core.message.ToolUseBlock; @@ -51,7 +51,7 @@ *
  • Replaced all {@code legacy.hook.*} events with {@code MiddlewareBase} override methods.
  • *
  • {@code PreCallEvent} / {@code PostCallEvent} → {@code onAgent()} before/after {@code next.apply()}.
  • *
  • {@code PreActingEvent} / {@code PostActingEvent} → {@code onActing()} before/after {@code next.apply()}.
  • - *
  • {@code ActingChunkEvent} → tap {@code ToolResultTextDeltaEvent} inside the {@code onActing()} stream.
  • + *
  • {@code ActingChunkEvent} → tap {@code ToolProgressEvent} inside the {@code onActing()} stream.
  • *
  • Reasoning streaming is observable via {@code onReasoning()} stream events.
  • *
  • {@code JsonlTraceExporter} removed; replaced by custom file-logging middleware pattern.
  • *
  • {@code .hooks(List)} → {@code .middlewares(List)}.
  • @@ -185,10 +185,10 @@ public Flux onActing( return next.apply(input) .doOnNext( event -> { - if (event instanceof ToolResultTextDeltaEvent delta) { + if (event instanceof ToolProgressEvent progress) { System.out.println( "[MIDDLEWARE] tool progress chunk: " - + delta.getDelta()); + + progress.getText()); } }) .doOnComplete( diff --git a/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/streaming/ToolStreamingExample.java b/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/streaming/ToolStreamingExample.java index de9d38d597..80ae026e68 100644 --- a/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/streaming/ToolStreamingExample.java +++ b/agentscope-examples/documentation/src/main/java/io/agentscope/examples/documentation2/streaming/ToolStreamingExample.java @@ -17,11 +17,12 @@ import io.agentscope.core.ReActAgent; import io.agentscope.core.event.AgentEvent; -import io.agentscope.core.event.ToolResultDataDeltaEvent; +import io.agentscope.core.event.ToolProgressEvent; import io.agentscope.core.event.ToolResultEndEvent; import io.agentscope.core.event.ToolResultStartEvent; import io.agentscope.core.event.ToolResultTextDeltaEvent; import io.agentscope.core.message.Msg; +import io.agentscope.core.message.TextBlock; import io.agentscope.core.message.ToolResultBlock; import io.agentscope.core.message.UserMessage; import io.agentscope.core.tool.Tool; @@ -31,8 +32,7 @@ /** * ToolStreamingExample - Demonstrates how tools emit streaming progress via {@link ToolEmitter} - * and how the caller receives them as {@link ToolResultTextDeltaEvent} / {@link - * ToolResultDataDeltaEvent} through {@code streamEvents()}. + * and how the caller receives them as {@link ToolProgressEvent} through {@code streamEvents()}. * *

    How tool streaming works: * @@ -41,19 +41,19 @@ * the framework auto-injects it). *

  • During execution, the tool calls {@code emitter.emit(ToolResultBlock.text(...))} to push * intermediate progress chunks. - *
  • The framework routes each chunk to the event stream as a {@link ToolResultTextDeltaEvent} - * or {@link ToolResultDataDeltaEvent}. + *
  • The framework routes each chunk to the event stream as a {@link ToolProgressEvent}. *
  • The tool's final {@code return} value becomes the tool result sent to the LLM — emitted - * chunks are NOT sent to the LLM (they are for the UI only). + * chunks are NOT sent to the LLM (they are for the UI only), and they arrive separately as + * {@link ToolResultTextDeltaEvent}. * * *

    Event sequence for a streaming tool call: *

      *   TOOL_RESULT_START        (tool execution begins)
    - *     TOOL_RESULT_TEXT_DELTA  ("Step 1: Fetching data...")
    - *     TOOL_RESULT_TEXT_DELTA  ("Step 2: Processing...")
    - *     TOOL_RESULT_TEXT_DELTA  ("Step 3: Formatting results...")
    - *   TOOL_RESULT_END          (state=SUCCESS, final result sent to LLM)
    + *     TOOL_PROGRESS           ("Step 1: Fetching data...")   — UI only, emitted by the tool
    + *     TOOL_PROGRESS           ("Step 2: Processing...")
    + *     TOOL_RESULT_TEXT_DELTA  ("Quantum computing...")       — the return value, sent to the LLM
    + *   TOOL_RESULT_END          (state=SUCCESS, result repeated on finalResultText)
      * 
    * *

    Run: @@ -101,17 +101,23 @@ private static void handleEvent(AgentEvent event) { System.out.printf( "%n[Tool Started] %s (id=%s)%n", e.getToolCallName(), e.getToolCallId()); - } else if (event instanceof ToolResultTextDeltaEvent e) { - // Streaming text chunks from the tool — display them incrementally. - // These come from ToolEmitter.emit() calls inside the tool method. - System.out.println(" ▸ " + e.getDelta()); + } else if (event instanceof ToolProgressEvent e) { + // Progress chunks from the tool's emitter.emit() calls — UI only, never the result. + if (e.getContent() instanceof TextBlock text) { + System.out.println(" ▸ " + text.getText()); + } else { + System.out.printf( + " ▸ [data block: %s]%n", e.getContent().getClass().getSimpleName()); + } - } else if (event instanceof ToolResultDataDeltaEvent e) { - // Non-text data blocks (images, binary, structured data). - System.out.printf(" ▸ [data block: %s]%n", e.getData().getClass().getSimpleName()); + } else if (event instanceof ToolResultTextDeltaEvent e) { + // The tool's return value, which is also what the model received. + System.out.println(" = " + e.getDelta()); } else if (event instanceof ToolResultEndEvent e) { - System.out.printf("[Tool Finished] id=%s state=%s%n", e.getToolCallId(), e.getState()); + System.out.printf( + "[Tool Finished] id=%s state=%s result=%s%n", + e.getToolCallId(), e.getState(), e.getFinalResultText()); } } diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-a2a/agentscope-extensions-a2a-server/src/main/java/io/agentscope/core/a2a/server/executor/AgentScopeAgentExecutor.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-a2a/agentscope-extensions-a2a-server/src/main/java/io/agentscope/core/a2a/server/executor/AgentScopeAgentExecutor.java index c4909176a3..c4223f59c2 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-a2a/agentscope-extensions-a2a-server/src/main/java/io/agentscope/core/a2a/server/executor/AgentScopeAgentExecutor.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-a2a/agentscope-extensions-a2a-server/src/main/java/io/agentscope/core/a2a/server/executor/AgentScopeAgentExecutor.java @@ -40,6 +40,7 @@ import io.agentscope.core.event.HintBlockEvent; import io.agentscope.core.event.TextBlockDeltaEvent; import io.agentscope.core.event.ThinkingBlockDeltaEvent; +import io.agentscope.core.event.ToolProgressEvent; import io.agentscope.core.event.ToolResultDataDeltaEvent; import io.agentscope.core.event.ToolResultTextDeltaEvent; import io.agentscope.core.message.ContentBlock; @@ -276,6 +277,7 @@ private Set generateRequiredEventTypes( return Set.of( AgentEventType.TEXT_BLOCK_DELTA, AgentEventType.THINKING_BLOCK_DELTA, + AgentEventType.TOOL_PROGRESS, AgentEventType.TOOL_RESULT_TEXT_DELTA, AgentEventType.TOOL_RESULT_DATA_DELTA, AgentEventType.HINT_BLOCK); @@ -343,6 +345,21 @@ private Msg convertToResponseMessage(AgentEvent output) { event.getReplyId(), ThinkingBlock.builder().thinking(event.getDelta()).build()); } + if (output instanceof ToolProgressEvent event) { + ContentBlock progress = event.getContent(); + if (progress == null) { + return null; + } + return toolResultMessage( + event, + event.getReplyId(), + ToolResultBlock.builder() + .id(event.getToolCallId()) + .name(event.getToolCallName()) + .output(progress) + .metadata(event.getMetadata()) + .build()); + } if (output instanceof ToolResultTextDeltaEvent event) { return toolResultMessage( event, @@ -523,6 +540,7 @@ protected void handleEvent(AgentEvent output, Msg responseMessage) { private boolean isStreamingChunk(AgentEvent output) { return output instanceof TextBlockDeltaEvent || output instanceof ThinkingBlockDeltaEvent + || output instanceof ToolProgressEvent || output instanceof ToolResultTextDeltaEvent || output instanceof ToolResultDataDeltaEvent; } diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-a2a/agentscope-extensions-a2a-server/src/test/java/io/agentscope/core/a2a/server/executor/AgentScopeAgentExecutorTest.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-a2a/agentscope-extensions-a2a-server/src/test/java/io/agentscope/core/a2a/server/executor/AgentScopeAgentExecutorTest.java index 0b3fa5a1e1..34541e8de8 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-a2a/agentscope-extensions-a2a-server/src/test/java/io/agentscope/core/a2a/server/executor/AgentScopeAgentExecutorTest.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-a2a/agentscope-extensions-a2a-server/src/test/java/io/agentscope/core/a2a/server/executor/AgentScopeAgentExecutorTest.java @@ -55,8 +55,10 @@ import io.agentscope.core.event.AgentStartEvent; import io.agentscope.core.event.HintBlockEvent; import io.agentscope.core.event.TextBlockDeltaEvent; +import io.agentscope.core.event.ToolProgressEvent; import io.agentscope.core.event.ToolResultTextDeltaEvent; import io.agentscope.core.message.Msg; +import io.agentscope.core.message.TextBlock; import java.time.Duration; import java.util.LinkedList; import java.util.List; @@ -429,6 +431,68 @@ void testExecuteAgentWithStreamingRequestPreservesEventMetadata() throws JSONRPC assertEquals("trace-1", toolPart.getMetadata().get("traceId")); } + @Test + @DisplayName("Should forward tool progress while inner events are enabled") + void testToolProgressReachesClientWhenInnerEventsEnabled() throws JSONRPCError { + AgentExecuteProperties properties = + AgentExecuteProperties.builder().requireInnerMessage(true).build(); + executor = new AgentScopeAgentExecutor(mockAgentRunner, properties); + doMockForContext(true, false, false); + when(mockAgentRunner.streamEvents(anyList(), any(AgentRequestOptions.class))) + .thenReturn( + Flux.just( + new ToolProgressEvent( + "reply-id", + "tool-call-id", + "search", + TextBlock.builder().text("scanning 3 dirs").build()))); + AtomicReference> messageRef = mockStreamingEventQueueRef(); + + executor.execute(mockContext, mockEventQueue); + + List artifacts = + messageRef.get().stream() + .filter(TaskArtifactUpdateEvent.class::isInstance) + .map(TaskArtifactUpdateEvent.class::cast) + .map(TaskArtifactUpdateEvent::getArtifact) + .toList(); + assertEquals( + 1, + artifacts.size(), + "progress moved to its own event type but still reaches the client"); + DataPart part = assertInstanceOf(DataPart.class, artifacts.get(0).parts().get(0)); + assertEquals( + Boolean.TRUE, + part.getMetadata().get(MessageConstants.STREAM_CHUNK_METADATA_KEY)); + } + + @Test + @DisplayName("Should withhold tool progress when inner events are disabled") + void testToolProgressIsWithheldWhenInnerEventsDisabled() throws JSONRPCError { + doMockForContext(true, false, false); + when(mockAgentRunner.streamEvents(anyList(), any(AgentRequestOptions.class))) + .thenReturn( + Flux.just( + new ToolProgressEvent( + "reply-id", + "tool-call-id", + "search", + TextBlock.builder().text("scanning 3 dirs").build()))); + AtomicReference> messageRef = mockStreamingEventQueueRef(); + + executor.execute(mockContext, mockEventQueue); + + List artifacts = + messageRef.get().stream() + .filter(TaskArtifactUpdateEvent.class::isInstance) + .map(TaskArtifactUpdateEvent.class::cast) + .map(TaskArtifactUpdateEvent::getArtifact) + .toList(); + assertTrue( + artifacts.isEmpty(), + "the default configuration must not start streaming progress: " + artifacts); + } + @Test @DisplayName( "Should execute agent and process streaming request with inner event but disabled") 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 a646489683..336faf4d28 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,7 @@ public class AguiStreamContext { private String currentTextMessageId; private String currentReasoningMessageId; private final Map toolResultContent = new LinkedHashMap<>(); + private final Map pendingInterrupts = new LinkedHashMap<>(); private final Set warnedMissingToolCallIdOperations = new LinkedHashSet<>(); private final TokenUsageAccumulator tokenUsageAccumulator = new TokenUsageAccumulator(); @@ -282,21 +283,36 @@ public void appendToolResultData(String toolCallId, ContentBlock data) { } public void endToolResult(String replyId, String toolCallId) { + endToolResult(replyId, toolCallId, null); + } + + /** + * Close out a tool call and emit its {@code TOOL_CALL_RESULT}. + * + *

    The buffered deltas are the result: {@code ToolResult*DeltaEvent} carries the tool method's + * return value, while progress published through {@code ToolEmitter} travels as {@code + * ToolProgressEvent}, which this context does not buffer at all. {@code finalResultText} as + * reported by {@link io.agentscope.core.event.ToolResultEndEvent} is the fallback for a result + * that produced no bufferable delta, and that is also what separates a blank return value + * ({@code ""}) from a producer that reported none ({@code null}). + * + * @param replyId the enclosing reply id + * @param toolCallId the tool call being closed + * @param finalResultText the return value, used when nothing was buffered + */ + public void endToolResult(String replyId, String toolCallId, String finalResultText) { if (!hasKnownToolCall(toolCallId, "ToolResultEndEvent")) { return; } if (startedToolCalls.contains(toolCallId) && endedToolCalls.add(toolCallId)) { emit(new AguiEvent.ToolCallEnd(threadId, runId, toolCallId)); } - StringBuilder content = toolResultContent.remove(toolCallId); + StringBuilder buffered = toolResultContent.remove(toolCallId); + String content = + buffered != null && !buffered.isEmpty() ? buffered.toString() : finalResultText; emit( new AguiEvent.ToolCallResult( - threadId, - runId, - toolCallId, - content != null && !content.isEmpty() ? content.toString() : null, - "tool", - replyId + ":" + toolCallId)); + threadId, runId, toolCallId, content, "tool", replyId + ":" + toolCallId)); } public void markToolCallSuspended(String toolCallId) { diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ToolResultEventConverter.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ToolResultEventConverter.java index 2be5ab85ae..a5f421b129 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ToolResultEventConverter.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/strategy/ToolResultEventConverter.java @@ -16,6 +16,7 @@ package io.agentscope.core.agui.adapter.strategy; import io.agentscope.core.event.AgentEvent; +import io.agentscope.core.event.ToolProgressEvent; import io.agentscope.core.event.ToolResultDataDeltaEvent; import io.agentscope.core.event.ToolResultEndEvent; import io.agentscope.core.event.ToolResultStartEvent; @@ -29,6 +30,7 @@ final class ToolResultEventConverter implements AgentEventConverter { public Set> eventTypes() { return Set.of( ToolResultStartEvent.class, + ToolProgressEvent.class, ToolResultTextDeltaEvent.class, ToolResultDataDeltaEvent.class, ToolResultEndEvent.class); @@ -36,6 +38,13 @@ public Set> eventTypes() { @Override public void convert(AgentEvent event, AguiStreamContext context) { + if (event instanceof ToolProgressEvent) { + // AG-UI has no channel for tool progress, and the result it renders comes from the + // deltas plus the authoritative return value. Buffering progress here is precisely the + // leak the result events now avoid by carrying their own type, so a progress chunk is + // deliberately dropped rather than passed through as a raw event. + return; + } if (event instanceof ToolResultTextDeltaEvent textDelta) { context.appendToolResultText(textDelta.getToolCallId(), textDelta.getDelta()); } else if (event instanceof ToolResultDataDeltaEvent dataDelta) { @@ -50,7 +59,7 @@ public void convert(AgentEvent event, AguiStreamContext context) { context.markToolCallSuspended(end.getToolCallId()); return; } - context.endToolResult(end.getReplyId(), end.getToolCallId()); + context.endToolResult(end.getReplyId(), end.getToolCallId(), end.getFinalResultText()); } } } 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 91f20654b0..d465ca86b5 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 @@ -61,6 +61,7 @@ import io.agentscope.core.event.ToolCallDeltaEvent; import io.agentscope.core.event.ToolCallEndEvent; import io.agentscope.core.event.ToolCallStartEvent; +import io.agentscope.core.event.ToolProgressEvent; import io.agentscope.core.event.ToolResultDataDeltaEvent; import io.agentscope.core.event.ToolResultEndEvent; import io.agentscope.core.event.ToolResultStartEvent; @@ -69,11 +70,13 @@ import io.agentscope.core.message.AssistantMessage; import io.agentscope.core.message.ContentBlock; import io.agentscope.core.message.GenerateReason; +import io.agentscope.core.message.ImageBlock; import io.agentscope.core.message.Msg; import io.agentscope.core.message.TextBlock; import io.agentscope.core.message.ToolResultBlock; import io.agentscope.core.message.ToolResultState; import io.agentscope.core.message.ToolUseBlock; +import io.agentscope.core.message.URLSource; import io.agentscope.core.model.ChatUsage; import io.agentscope.core.model.ToolSchema; import io.agentscope.core.tool.SchemaOnlyTool; @@ -755,6 +758,70 @@ void testToolResultDeltasAreAggregatedIntoToolCallResult() { assertEquals("reply-tool:tool-1", result.messageId()); } + @Test + void testProgressChunksBecomeNeitherTheResultNorRawEvents() { + // ToolEmitter progress travels as ToolProgressEvent now. AG-UI has no channel for it, + // so + // the converter drops it: TOOL_CALL_RESULT is what a client stores as the tool message, + // and a progress event that reached the registry fallback would instead surface as one + // RAW event per chunk. + List events = + runReActEvents( + new ToolCallStartEvent("reply-emitter", "tool-1", "enroll"), + new ToolCallEndEvent("reply-emitter", "tool-1", "enroll"), + new ToolResultStartEvent("reply-emitter", "tool-1", "enroll"), + progress("reply-emitter", "tool-1", "enroll", "press once"), + progress("reply-emitter", "tool-1", "enroll", "press again"), + new ToolResultTextDeltaEvent( + "reply-emitter", "tool-1", "enroll", "fingerprint enrolled"), + new ToolResultEndEvent( + "reply-emitter", + "tool-1", + "enroll", + null, + "fingerprint enrolled")); + + assertEquals( + "fingerprint enrolled", + firstToolCallResult(events).content(), + "the result delta is what the client stores as the tool message"); + assertTrue( + events.stream().noneMatch(AguiEvent.Raw.class::isInstance), + "progress is handled by the converter rather than passed through as a raw" + + " event"); + } + + @Test + void testToolResultWithoutFinalTextMetadataStillUsesDeltaBuffer() { + List events = + runReActEvents( + new ToolCallStartEvent("reply-legacy", "tool-1", "lookup"), + new ToolCallEndEvent("reply-legacy", "tool-1", "lookup"), + new ToolResultStartEvent("reply-legacy", "tool-1", "lookup"), + new ToolResultTextDeltaEvent( + "reply-legacy", "tool-1", "lookup", "partial"), + new ToolResultEndEvent("reply-legacy", "tool-1", "lookup", null)); + + assertEquals( + "partial", + firstToolCallResult(events).content(), + "producers that do not report a return value keep the previous behaviour"); + } + + private static AguiEvent.ToolCallResult firstToolCallResult(List events) { + return events.stream() + .filter(AguiEvent.ToolCallResult.class::isInstance) + .map(AguiEvent.ToolCallResult.class::cast) + .findFirst() + .orElseThrow(); + } + + private static ToolProgressEvent progress( + String replyId, String toolCallId, String toolCallName, String text) { + return new ToolProgressEvent( + replyId, toolCallId, toolCallName, TextBlock.builder().text(text).build()); + } + @Test void testParallelToolResultsUsePerToolMessageIds() { List messageIds = @@ -824,6 +891,103 @@ void testToolResultTextAndDataDeltasAreJoinedInArrivalOrder() { assertEquals("hel\nstructuredlo", result.content()); } + @Test + void testMultimodalResultKeepsEveryPartInArrivalOrder() { + // A producer emits each block of the return value as its own delta, so the buffer holds + // the whole result in the order the parts arrived, ahead of the end event's joined + // text. + ImageBlock image = + ImageBlock.builder() + .source(URLSource.builder().url("https://example.com/cat.png").build()) + .build(); + List events = + runReActEvents( + new ToolCallStartEvent("reply-tool", "tool-1", "lookup"), + new ToolCallEndEvent("reply-tool", "tool-1", "lookup"), + new ToolResultStartEvent("reply-tool", "tool-1", "lookup"), + new ToolResultDataDeltaEvent("reply-tool", "tool-1", "lookup", image), + new ToolResultTextDeltaEvent( + "reply-tool", "tool-1", "lookup", "3 rows found"), + new ToolResultEndEvent( + "reply-tool", "tool-1", "lookup", null, "3 rows found")); + + String content = firstToolCallResult(events).content(); + int imageAt = content.indexOf("https://example.com/cat.png"); + int textAt = content.indexOf("3 rows found"); + assertTrue(imageAt >= 0, "the image part is kept: " + content); + assertTrue(textAt >= 0, "the text part is kept: " + content); + assertTrue( + imageAt < textAt, + "the buffered deltas win, so the parts stay in arrival order: " + content); + } + + @Test + void testResultTextFallsBackToTheEndEventWhenNothingWasBuffered() { + // A producer that reports no deltas still has to surface its return value, and "" would + // read as an empty result rather than as nothing reported. + List events = + runReActEvents( + new ToolCallStartEvent("reply-empty", "tool-1", "lookup"), + new ToolCallEndEvent("reply-empty", "tool-1", "lookup"), + new ToolResultStartEvent("reply-empty", "tool-1", "lookup"), + new ToolResultEndEvent( + "reply-empty", + "tool-1", + "lookup", + null, + "reported on the end event")); + + assertEquals( + "reported on the end event", + firstToolCallResult(events).content(), + "the end event is the fallback when the delta stream carried nothing"); + } + + @Test + void testImageOnlyResultDoesNotFallBackToProgressText() { + // A tool that streams progress and returns only an image has no text in its return + // value, so the field is "" — reported, but empty. Reading that as "nothing reported" + // used to put the progress chunks back into TOOL_CALL_RESULT. + ImageBlock image = + ImageBlock.builder() + .source( + URLSource.builder() + .url("https://example.com/chart.png") + .build()) + .build(); + List events = + runReActEvents( + new ToolCallStartEvent("reply-image", "tool-1", "render"), + new ToolCallEndEvent("reply-image", "tool-1", "render"), + new ToolResultStartEvent("reply-image", "tool-1", "render"), + progress("reply-image", "tool-1", "render", "drawing axes"), + new ToolResultDataDeltaEvent("reply-image", "tool-1", "render", image), + new ToolResultEndEvent("reply-image", "tool-1", "render", null, "")); + + String content = firstToolCallResult(events).content(); + assertFalse( + content.contains("drawing axes"), "progress text must not leak: " + content); + assertTrue( + content.contains("https://example.com/chart.png"), + "the image is the result: " + content); + } + + @Test + void testBlankReturnValueYieldsEmptyResultNotProgressText() { + List events = + runReActEvents( + new ToolCallStartEvent("reply-blank", "tool-1", "cleanup"), + new ToolCallEndEvent("reply-blank", "tool-1", "cleanup"), + new ToolResultStartEvent("reply-blank", "tool-1", "cleanup"), + progress("reply-blank", "tool-1", "cleanup", "removing 3 files"), + new ToolResultEndEvent("reply-blank", "tool-1", "cleanup", null, "")); + + assertEquals( + "", + firstToolCallResult(events).content(), + "an empty return value reads as empty, not as what the tool streamed"); + } + @Test void testToolResultEndWithoutContentStillEmitsNullResult() { List events = diff --git a/agentscope-harness/src/test/java/io/agentscope/harness/agent/subagent/protocol/RemoteEventCodecPassthroughTest.java b/agentscope-harness/src/test/java/io/agentscope/harness/agent/subagent/protocol/RemoteEventCodecPassthroughTest.java index 40980429e3..1438748ebb 100644 --- a/agentscope-harness/src/test/java/io/agentscope/harness/agent/subagent/protocol/RemoteEventCodecPassthroughTest.java +++ b/agentscope-harness/src/test/java/io/agentscope/harness/agent/subagent/protocol/RemoteEventCodecPassthroughTest.java @@ -50,6 +50,7 @@ import io.agentscope.core.event.ToolCallDeltaEvent; import io.agentscope.core.event.ToolCallEndEvent; import io.agentscope.core.event.ToolCallStartEvent; +import io.agentscope.core.event.ToolProgressEvent; import io.agentscope.core.event.ToolResultDataDeltaEvent; import io.agentscope.core.event.ToolResultEndEvent; import io.agentscope.core.event.ToolResultStartEvent; @@ -126,6 +127,13 @@ private static Map sampleEvents() { events.put( AgentEventType.TOOL_RESULT_START, new ToolResultStartEvent("reply", "call-1", "read_file")); + events.put( + AgentEventType.TOOL_PROGRESS, + new ToolProgressEvent( + "reply", + "call-1", + "read_file", + TextBlock.builder().text("opened a").build())); events.put( AgentEventType.TOOL_RESULT_TEXT_DELTA, new ToolResultTextDeltaEvent("reply", "call-1", "read_file", "line one"));