From cd9a4ed303a0b800c0e6cf40cc5b8da83672a2a2 Mon Sep 17 00:00:00 2001 From: xzxiaoshan <365384722@qq.com> Date: Mon, 3 Aug 2026 13:57:01 +0800 Subject: [PATCH 1/4] feat(context-enhance): add RuntimeContext onStateLoaded callback for anchor-based AG-UI message auto-merge --- .../java/io/agentscope/core/ReActAgent.java | 6 + .../agentscope/core/agent/RuntimeContext.java | 44 ++++ .../core/agent/RuntimeContextTest.java | 93 ++++++++ .../core/agui/adapter/AguiAgentAdapter.java | 32 +++ .../agui/processor/AguiRequestProcessor.java | 57 +---- .../AguiAgentAdapterMessageMergeTest.java | 225 ++++++++++++++++++ .../processor/AguiRequestProcessorTest.java | 37 ++- 7 files changed, 420 insertions(+), 74 deletions(-) create mode 100644 agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterMessageMergeTest.java 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 562f1f9ed6..5df1a53132 100644 --- a/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java +++ b/agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java @@ -138,6 +138,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.BiConsumer; import java.util.function.Function; import java.util.stream.Collectors; import java.util.stream.IntStream; @@ -546,6 +547,11 @@ protected Object beforeAgentExecution(List msgs, RuntimeContext rc) { // the active session's state via rc.getAgentState() (call-scoped, concurrency-safe) // rather than agent.getAgentState() (not call-scoped under concurrency). ctx.setAgentState(scope.state); + // state 加载完成,触发 onStateLoaded 回调 + BiConsumer> onStateLoaded = ctx.getOnStateLoaded(); + if (onStateLoaded != null) { + onStateLoaded.accept(ctx, msgs); + } this.activeRc = ctx; bindRuntimeContextToHooks(ctx); // Seed per-call state onto the active execution scope. The system message is initialised diff --git a/agentscope-core/src/main/java/io/agentscope/core/agent/RuntimeContext.java b/agentscope-core/src/main/java/io/agentscope/core/agent/RuntimeContext.java index 69c1a71f7d..54fac14b90 100644 --- a/agentscope-core/src/main/java/io/agentscope/core/agent/RuntimeContext.java +++ b/agentscope-core/src/main/java/io/agentscope/core/agent/RuntimeContext.java @@ -15,13 +15,16 @@ */ package io.agentscope.core.agent; +import io.agentscope.core.message.Msg; import io.agentscope.core.state.AgentState; import io.agentscope.core.tool.ContextStore; import io.agentscope.core.tool.ToolExecutionContext; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.function.BiConsumer; /** * Per-call metadata for an agent run: session-scoped fields plus a thread-safe attribute bag and @@ -45,6 +48,14 @@ public class RuntimeContext { */ private volatile AgentState agentState; + /** + * Callback fired after {@link #agentState} is loaded and set on this context (inside the + * agent's {@code beforeAgentExecution}, right after {@link #setAgentState(AgentState)}). The + * callback receives this RuntimeContext and the mutable incoming message list — it may modify + * either in place. {@code null} when no callback is registered. + */ + private volatile BiConsumer> onStateLoaded; + /** String-keyed extras (legacy and generic extension). */ private final ConcurrentMap stringAttributes; @@ -63,6 +74,7 @@ private RuntimeContext(Builder builder) { this.typedAttributes = new ConcurrentHashMap<>(); this.toolExecutionContext = builder.toolExecutionContext; this.agentState = builder.agentState; + this.onStateLoaded = builder.onStateLoaded; if (builder.stringExtras != null) { this.stringAttributes.putAll(builder.stringExtras); } @@ -115,6 +127,25 @@ public void setAgentState(AgentState agentState) { this.agentState = agentState; } + /** + * Returns the callback fired after AgentState is loaded and set on this context, or {@code null} + * if none is registered. + */ + public BiConsumer> getOnStateLoaded() { + return onStateLoaded; + } + + /** + * Registers a callback fired after AgentState is loaded (inside {@code beforeAgentExecution}, + * right after {@link #setAgentState(AgentState)}). The callback receives this RuntimeContext and + * the mutable incoming message list — it may modify either in place. + * + * @param onStateLoaded the callback, or {@code null} to clear + */ + public void setOnStateLoaded(BiConsumer> onStateLoaded) { + this.onStateLoaded = onStateLoaded; + } + /** * Resolves the live {@link AgentState} for the current call, preferring the call-scoped state * carried on {@code ctx} (concurrency-safe) and falling back to {@code fallbackAgent}'s state @@ -329,6 +360,7 @@ public static class Builder { private final Map, Map> typedValues = new HashMap<>(); private ToolExecutionContext toolExecutionContext; private AgentState agentState; + private BiConsumer> onStateLoaded; public Builder sessionId(String sessionId) { this.sessionId = sessionId; @@ -345,6 +377,17 @@ public Builder agentState(AgentState agentState) { return this; } + /** + * Registers the {@link RuntimeContext#getOnStateLoaded()} callback on the built context. + * + * @param onStateLoaded the callback, or {@code null} to leave unset + * @return this builder + */ + public Builder onStateLoaded(BiConsumer> onStateLoaded) { + this.onStateLoaded = onStateLoaded; + return this; + } + public Builder put(String key, Object value) { if (this.stringExtras == null) { this.stringExtras = new ConcurrentHashMap<>(); @@ -383,6 +426,7 @@ public Builder from(RuntimeContext source) { this.sessionId = source.sessionId; this.userId = source.userId; this.agentState = source.agentState; + this.onStateLoaded = source.onStateLoaded; this.toolExecutionContext = source.toolExecutionContext; if (!source.stringAttributes.isEmpty()) { this.stringExtras = new ConcurrentHashMap<>(source.stringAttributes); diff --git a/agentscope-core/src/test/java/io/agentscope/core/agent/RuntimeContextTest.java b/agentscope-core/src/test/java/io/agentscope/core/agent/RuntimeContextTest.java index fb23b4b41f..e0c7f8cc1b 100644 --- a/agentscope-core/src/test/java/io/agentscope/core/agent/RuntimeContextTest.java +++ b/agentscope-core/src/test/java/io/agentscope/core/agent/RuntimeContextTest.java @@ -19,10 +19,18 @@ import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertTrue; +import io.agentscope.core.message.Msg; +import io.agentscope.core.message.UserMessage; +import io.agentscope.core.state.AgentState; import io.agentscope.core.tool.ToolExecutionContext; +import java.util.ArrayList; +import java.util.List; import java.util.concurrent.CyclicBarrier; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.BiConsumer; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; @@ -157,6 +165,91 @@ void builderCopyHandlesNullSource() { assertNull(empty.get("missing", Marker.class)); } + @Test + @DisplayName("onStateLoaded defaults to null") + void onStateLoaded_defaultsToNull() { + RuntimeContext ctx = RuntimeContext.empty(); + assertNull(ctx.getOnStateLoaded()); + } + + @Test + @DisplayName("builder.onStateLoaded propagates the callback to the built context") + void onStateLoaded_builderPropagatesCallback() { + BiConsumer> callback = (c, m) -> {}; + RuntimeContext ctx = RuntimeContext.builder().onStateLoaded(callback).build(); + assertSame(callback, ctx.getOnStateLoaded()); + } + + @Test + @DisplayName("setOnStateLoaded replaces and clears the callback") + void onStateLoaded_setterReplacesAndClears() { + BiConsumer> first = (c, m) -> {}; + BiConsumer> second = (c, m) -> {}; + RuntimeContext ctx = RuntimeContext.builder().onStateLoaded(first).build(); + assertSame(first, ctx.getOnStateLoaded()); + + ctx.setOnStateLoaded(second); + assertSame(second, ctx.getOnStateLoaded()); + + ctx.setOnStateLoaded(null); + assertNull(ctx.getOnStateLoaded()); + } + + @Test + @DisplayName("builder(source) preserves the onStateLoaded callback") + void onStateLoaded_builderCopyPreservesCallback() { + BiConsumer> callback = (c, m) -> {}; + RuntimeContext source = RuntimeContext.builder().onStateLoaded(callback).build(); + RuntimeContext copy = RuntimeContext.builder(source).build(); + assertSame(callback, copy.getOnStateLoaded()); + } + + @Test + @DisplayName("onStateLoaded callback receives the context and mutable message list") + void onStateLoaded_callbackReceivesContextAndMsgs() { + RuntimeContext ctx = RuntimeContext.empty(); + AtomicInteger fired = new AtomicInteger(); + ctx.setOnStateLoaded( + (c, m) -> { + fired.incrementAndGet(); + assertSame(c, ctx); + m.clear(); + }); + AgentState state = AgentState.builder().build(); + List msgs = new ArrayList<>(List.of(new UserMessage("hi"))); + + ctx.setAgentState(state); + ctx.getOnStateLoaded().accept(ctx, msgs); + + assertEquals(1, fired.get()); + assertTrue(msgs.isEmpty()); + } + + @Test + @DisplayName("onStateLoaded callback can read the agent state just set on the context") + void onStateLoaded_callbackReadsFreshlySetState() { + AgentState state = AgentState.builder().summary("loaded").build(); + AtomicReference seen = new AtomicReference<>(); + BiConsumer> callback = (c, m) -> seen.set(c.getAgentState()); + RuntimeContext ctx = RuntimeContext.builder().onStateLoaded(callback).build(); + ctx.setAgentState(state); + ctx.getOnStateLoaded().accept(ctx, List.of()); + + assertSame(state, seen.get()); + } + + @Test + @DisplayName("onStateLoaded can be registered on a copied context without affecting the source") + void onStateLoaded_copyIsIndependentAfterRegistration() { + BiConsumer> original = (c, m) -> {}; + RuntimeContext source = RuntimeContext.builder().onStateLoaded(original).build(); + BiConsumer> override = (c, m) -> {}; + RuntimeContext copy = RuntimeContext.builder(source).onStateLoaded(override).build(); + + assertSame(override, copy.getOnStateLoaded()); + assertSame(original, source.getOnStateLoaded()); + } + @Test @DisplayName("concurrent puts on distinct keys from multiple threads") void threadSafety() throws Exception { diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java index bc02850c61..e389074e16 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java @@ -36,6 +36,7 @@ import io.agentscope.core.message.ToolResultBlock; import io.agentscope.core.message.ToolUseBlock; import io.agentscope.core.model.ToolSchema; +import io.agentscope.core.state.AgentState; import io.agentscope.core.tool.AgentTool; import io.agentscope.core.tool.SchemaOnlyTool; import io.agentscope.core.tool.Toolkit; @@ -50,6 +51,7 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.BiConsumer; import java.util.function.Supplier; import reactor.core.publisher.Flux; @@ -306,9 +308,39 @@ protected RuntimeContext buildRuntimeContext( .put(RUNTIME_CONTEXT_STATE_KEY, input.getState()) .put(RUNTIME_CONTEXT_FORWARDED_PROPS_KEY, input.getForwardedProps()) .put(RUNTIME_CONTEXT_RESUME_KEY, input.getResume()) + .onStateLoaded(createMessageMergeHandler()) .build(); } + /** + * Creates the {@code onStateLoaded} callback that deduplicates incoming messages against the + * already-persisted AgentState context. + * + *

The last message id in {@code state.getContext()} is used as the anchor: if it is found in + * the incoming {@code msgs} list, the anchor and everything before it is removed in place (the + * input was a full transcript); otherwise the incoming list is treated as purely incremental and + * left untouched. When the context is empty or the state is null the callback is a no-op. + */ + private BiConsumer> createMessageMergeHandler() { + return (ctx, msgs) -> { + AgentState state = ctx.getAgentState(); + if (state == null) { + return; + } + List context = state.getContext(); + if (context.isEmpty()) { + return; + } + String anchorId = context.get(context.size() - 1).getId(); + for (int i = msgs.size() - 1; i >= 0; i--) { + if (anchorId.equals(msgs.get(i).getId())) { + msgs.subList(0, i + 1).clear(); + return; + } + } + }; + } + @SuppressWarnings("unchecked") private Map resumeToolCallIds(RuntimeContext runtimeContext) { if (runtimeContext == null) { diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/processor/AguiRequestProcessor.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/processor/AguiRequestProcessor.java index 9574d48eee..584c9ea943 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/processor/AguiRequestProcessor.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/processor/AguiRequestProcessor.java @@ -21,9 +21,7 @@ import io.agentscope.core.agui.adapter.AguiAgentAdapter; import io.agentscope.core.agui.adapter.AguiAgentAdapterFactory; import io.agentscope.core.agui.event.AguiEvent; -import io.agentscope.core.agui.model.AguiMessage; import io.agentscope.core.agui.model.RunAgentInput; -import java.util.List; import java.util.Objects; import java.util.concurrent.atomic.AtomicBoolean; import org.slf4j.Logger; @@ -39,7 +37,6 @@ *

Responsibilities: *

    *
  • Agent ID resolution from multiple sources
  • - *
  • Message extraction for server-side memory scenarios
  • *
  • Agent resolution via {@link AgentResolver}
  • *
  • Event stream generation via {@link AguiAgentAdapter}
  • *
@@ -136,16 +133,10 @@ public ProcessResult process( } try { - // Determine effective input based on server-side memory + // Full input is forwarded; message dedup against persisted + // AgentState context is handled by the onStateLoaded callback + // registered in AguiAgentAdapter.buildRuntimeContext(). RunAgentInput effectiveInput = input; - if (agentResolver.hasMemory(threadId)) { - logger.debug( - "Using server-side memory for thread {}, extracting" - + " latest user message", - threadId); - effectiveInput = extractLatestUserMessage(input); - } - RuntimeContext effectiveRuntimeContext = resumeCoordinator.addResumeToolCallIds( input, runtimeContext); @@ -257,48 +248,6 @@ public String resolveAgentId(RunAgentInput input, String headerAgentId, String p return "default"; } - /** - * Extract only the latest user message from the input. - * - *

This is used when server-side memory is enabled and the agent already - * has conversation history. Only the latest user message needs to be passed. - * - * @param input The original input - * @return A new input with only the latest user message - */ - public RunAgentInput extractLatestUserMessage(RunAgentInput input) { - List messages = input.getMessages(); - if (messages == null || messages.isEmpty()) { - return input; - } - - // Find the last user message - AguiMessage lastUserMessage = null; - for (int i = messages.size() - 1; i >= 0; i--) { - AguiMessage msg = messages.get(i); - if ("user".equalsIgnoreCase(msg.getRole())) { - lastUserMessage = msg; - break; - } - } - - if (lastUserMessage == null) { - return input; - } - - // Create new input with only the last user message - return RunAgentInput.builder() - .threadId(input.getThreadId()) - .runId(input.getRunId()) - .messages(List.of(lastUserMessage)) - .tools(input.getTools()) - .context(input.getContext()) - .state(input.getState()) - .forwardedProps(input.getForwardedProps()) - .resume(input.getResume()) - .build(); - } - /** * Creates a new builder for AguiRequestProcessor. * diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterMessageMergeTest.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterMessageMergeTest.java new file mode 100644 index 0000000000..17aef1e2e5 --- /dev/null +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/adapter/AguiAgentAdapterMessageMergeTest.java @@ -0,0 +1,225 @@ +/* + * 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.agui.adapter; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.mockito.Mockito.mock; + +import io.agentscope.core.agent.Agent; +import io.agentscope.core.agent.RuntimeContext; +import io.agentscope.core.agui.model.RunAgentInput; +import io.agentscope.core.message.Msg; +import io.agentscope.core.message.UserMessage; +import io.agentscope.core.state.AgentState; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.BiConsumer; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +/** + * Tests for the {@code onStateLoaded} message-merge handler registered by {@link + * AguiAgentAdapter#buildRuntimeContext}. + * + *

These tests focus on the merge behavior triggered after AgentState is loaded: when the incoming + * message list is a full transcript whose tail overlaps the persisted context, the overlap prefix is + * stripped; when the list is already incremental (no overlap), it is left untouched. + * + *

Assertions compare message ids (not Msg instances) since {@link Msg} uses identity equality. + */ +@DisplayName("AguiAgentAdapter onStateLoaded message merge") +class AguiAgentAdapterMessageMergeTest { + + @Test + @DisplayName("buildRuntimeContext registers a non-null onStateLoaded callback") + void buildRuntimeContext_registersCallback() { + RuntimeContext ctx = newContextWithState(null); + + assertNotNull(ctx.getOnStateLoaded()); + } + + @Test + @DisplayName("callback receives the exact context it was registered on") + void merge_callbackReceivesSameContext() { + AgentState state = AgentState.builder().context(List.of(msg("m1"))).build(); + RuntimeContext ctx = newContextWithState(state); + AtomicReference seen = new AtomicReference<>(); + BiConsumer> original = ctx.getOnStateLoaded(); + ctx.setOnStateLoaded( + (c, m) -> { + seen.set(c); + original.accept(c, m); + }); + + fireCallback(ctx, new ArrayList<>(List.of(msg("m1")))); + + assertSame(ctx, seen.get()); + } + + @Test + @DisplayName("full transcript input: anchor hit strips the overlapping prefix") + void merge_anchorHit_stripsOverlappingPrefix() { + AgentState state = + AgentState.builder().context(List.of(msg("m1"), msg("m2"), msg("m3"))).build(); + RuntimeContext ctx = newContextWithState(state); + List incoming = + new ArrayList<>(List.of(msg("m1"), msg("m2"), msg("m3"), msg("m4"), msg("m5"))); + + fireCallback(ctx, incoming); + + assertEquals(List.of("m4", "m5"), idsOf(incoming)); + } + + @Test + @DisplayName("incremental input: anchor miss leaves msgs untouched") + void merge_anchorMiss_leavesMsgsUntouched() { + AgentState state = + AgentState.builder().context(List.of(msg("m1"), msg("m2"), msg("m3"))).build(); + RuntimeContext ctx = newContextWithState(state); + List incoming = new ArrayList<>(List.of(msg("m4"), msg("m5"))); + + fireCallback(ctx, incoming); + + assertEquals(List.of("m4", "m5"), idsOf(incoming)); + } + + @Test + @DisplayName("empty persisted context: no-op, msgs untouched") + void merge_emptyContext_leavesMsgsUntouched() { + AgentState state = AgentState.builder().build(); + RuntimeContext ctx = newContextWithState(state); + List incoming = new ArrayList<>(List.of(msg("m1"), msg("m2"), msg("m3"))); + + fireCallback(ctx, incoming); + + assertEquals(List.of("m1", "m2", "m3"), idsOf(incoming)); + } + + @Test + @DisplayName("null agent state on context: no-op, msgs untouched") + void merge_nullState_leavesMsgsUntouched() { + RuntimeContext ctx = newContextWithState(null); + List incoming = new ArrayList<>(List.of(msg("m1"), msg("m2"))); + + fireCallback(ctx, incoming); + + assertEquals(List.of("m1", "m2"), idsOf(incoming)); + } + + @Test + @DisplayName("empty input msgs: no-op even when context has anchor") + void merge_emptyIncomingMsgs_noOp() { + AgentState state = AgentState.builder().context(List.of(msg("m1"), msg("m2"))).build(); + RuntimeContext ctx = newContextWithState(state); + List incoming = new ArrayList<>(); + + fireCallback(ctx, incoming); + + assertEquals(List.of(), idsOf(incoming)); + } + + @Test + @DisplayName("anchor hit at position 0 keeps only messages after it") + void merge_anchorAtStart_clearsPrefix() { + AgentState state = AgentState.builder().context(List.of(msg("m1"))).build(); + RuntimeContext ctx = newContextWithState(state); + List incoming = new ArrayList<>(List.of(msg("m1"), msg("m2"))); + + fireCallback(ctx, incoming); + + assertEquals(List.of("m2"), idsOf(incoming)); + } + + @Test + @DisplayName("anchor is the last element of incoming: all duplicates removed") + void merge_anchorIsLastIncoming_removesAll() { + AgentState state = + AgentState.builder().context(List.of(msg("m1"), msg("m2"), msg("m3"))).build(); + RuntimeContext ctx = newContextWithState(state); + List incoming = new ArrayList<>(List.of(msg("m1"), msg("m2"), msg("m3"))); + + fireCallback(ctx, incoming); + + assertEquals(List.of(), idsOf(incoming)); + } + + @Test + @DisplayName("full and incremental inputs converge to the same effective msgs") + void merge_fullVsIncremental_converge() { + AgentState state = + AgentState.builder().context(List.of(msg("m1"), msg("m2"), msg("m3"))).build(); + + RuntimeContext ctxFull = newContextWithState(state); + List fullIncoming = + new ArrayList<>(List.of(msg("m1"), msg("m2"), msg("m3"), msg("m4"), msg("m5"))); + fireCallback(ctxFull, fullIncoming); + + RuntimeContext ctxInc = newContextWithState(state); + List incIncoming = new ArrayList<>(List.of(msg("m4"), msg("m5"))); + fireCallback(ctxInc, incIncoming); + + assertEquals(idsOf(incIncoming), idsOf(fullIncoming)); + } + + @Test + @DisplayName("firing the callback twice on the same already-merged list is stable") + void merge_firedTwice_isStable() { + AgentState state = + AgentState.builder().context(List.of(msg("m1"), msg("m2"), msg("m3"))).build(); + RuntimeContext ctx = newContextWithState(state); + List incoming = new ArrayList<>(List.of(msg("m1"), msg("m2"), msg("m3"), msg("m4"))); + + fireCallback(ctx, incoming); + List afterFirst = idsOf(incoming); + fireCallback(ctx, incoming); + + assertEquals(afterFirst, idsOf(incoming)); + } + + // ---------- helpers ---------- + + /** + * Builds a RuntimeContext via {@link AguiAgentAdapter#buildRuntimeContext}, then sets the given + * AgentState onto it — mirroring what {@code beforeAgentExecution} does right before firing + * {@code onStateLoaded}. + */ + private static RuntimeContext newContextWithState(AgentState state) { + AguiAgentAdapter adapter = + new AguiAgentAdapter(mock(Agent.class), AguiAdapterConfig.defaultConfig()); + RunAgentInput input = RunAgentInput.builder().threadId("t-1").runId("r-1").build(); + RuntimeContext ctx = adapter.buildRuntimeContext(input, null); + ctx.setAgentState(state); + return ctx; + } + + private static void fireCallback(RuntimeContext ctx, List msgs) { + BiConsumer> callback = ctx.getOnStateLoaded(); + assertNotNull(callback); + callback.accept(ctx, msgs); + } + + private static Msg msg(String id) { + Msg template = new UserMessage(id); + return Msg.builder().id(id).role(template.getRole()).content(template.getContent()).build(); + } + + private static List idsOf(List msgs) { + return msgs.stream().map(Msg::getId).toList(); + } +} diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/processor/AguiRequestProcessorTest.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/processor/AguiRequestProcessorTest.java index eed8b9bb1f..4d5a123d1b 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/processor/AguiRequestProcessorTest.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/test/java/io/agentscope/core/agui/processor/AguiRequestProcessorTest.java @@ -50,36 +50,34 @@ class AguiRequestProcessorTest { @Test - void extractLatestUserMessagePreservesFullRunInputMetadata() { + void processForwardsFullInputWithoutDroppingEarlierMessages() { + AgentResolver resolver = mock(AgentResolver.class); + ReActAgent agent = mock(ReActAgent.class); + ArgumentCaptor> msgsCaptor = ArgumentCaptor.forClass(List.class); + when(resolver.resolveAgent("default", "thread-1")).thenReturn(agent); + when(agent.streamEvents(msgsCaptor.capture(), any(RuntimeContext.class))) + .thenReturn(Flux.empty()); AguiRequestProcessor processor = - AguiRequestProcessor.builder().agentResolver(mock(AgentResolver.class)).build(); - AguiMessage firstUser = AguiMessage.userMessage("msg-1", "first"); - AguiMessage lastUser = AguiMessage.userMessage("msg-3", "last"); + AguiRequestProcessor.builder().agentResolver(resolver).build(); RunAgentInput input = RunAgentInput.builder() .threadId("thread-1") .runId("run-1") .messages( List.of( - firstUser, + AguiMessage.userMessage("msg-1", "first"), AguiMessage.assistantMessage("msg-2", "ok"), - lastUser)) - .state(Map.of("cursor", 8)) - .forwardedProps(Map.of("agentId", "agent-a")) - .resume( - List.of( - new AguiResume( - "int-1", - AguiResume.STATUS_RESOLVED, - Map.of("approved", true)))) + AguiMessage.userMessage("msg-3", "last"))) .build(); - RunAgentInput extracted = processor.extractLatestUserMessage(input); + processor.process(input, null, null).events().collectList().block(); - assertEquals(List.of(lastUser), extracted.getMessages()); - assertEquals(input.getState(), extracted.getState()); - assertEquals(input.getForwardedProps(), extracted.getForwardedProps()); - assertEquals(input.getResume(), extracted.getResume()); + // All three messages reach the agent — none are filtered out by the processor. + List forwarded = msgsCaptor.getValue(); + assertEquals(3, forwarded.size()); + assertEquals("msg-1", forwarded.get(0).getId()); + assertEquals("msg-2", forwarded.get(1).getId()); + assertEquals("msg-3", forwarded.get(2).getId()); } @Test @@ -121,7 +119,6 @@ void processRecordsInterruptsAndResolvesOfficialResumeToToolCallId() { AgentResolver resolver = mock(AgentResolver.class); ReActAgent agent = mock(ReActAgent.class); when(resolver.resolveAgent("default", "thread-1")).thenReturn(agent); - when(resolver.hasMemory("thread-1")).thenReturn(false); ArgumentCaptor> msgsCaptor = ArgumentCaptor.forClass(List.class); when(agent.streamEvents(msgsCaptor.capture(), any(RuntimeContext.class))) .thenReturn(Flux.just(new AgentEndEvent("reply-2"))); From 1d9674cec3c75a021f0e841ff2b2eee096a2854d Mon Sep 17 00:00:00 2001 From: xzxiaoshan <365384722@qq.com> Date: Mon, 3 Aug 2026 14:14:38 +0800 Subject: [PATCH 2/4] feat(context-enhance): The use of context provides a method for getLast to retrieve the last message --- .../java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java index e389074e16..a081ba9bd8 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java @@ -331,7 +331,7 @@ private BiConsumer> createMessageMergeHandler() { if (context.isEmpty()) { return; } - String anchorId = context.get(context.size() - 1).getId(); + String anchorId = context.getLast().getId(); for (int i = msgs.size() - 1; i >= 0; i--) { if (anchorId.equals(msgs.get(i).getId())) { msgs.subList(0, i + 1).clear(); From 24a84fd4b1b308f42bd99616761b21018684f4ac Mon Sep 17 00:00:00 2001 From: xzxiaoshan <365384722@qq.com> Date: Mon, 3 Aug 2026 14:47:25 +0800 Subject: [PATCH 3/4] fix(agui): replace List.getLast() with Java 17-compatible index access --- .../java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java index a081ba9bd8..e389074e16 100644 --- a/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java +++ b/agentscope-extensions/agentscope-extensions-protocol/agentscope-extensions-agui/src/main/java/io/agentscope/core/agui/adapter/AguiAgentAdapter.java @@ -331,7 +331,7 @@ private BiConsumer> createMessageMergeHandler() { if (context.isEmpty()) { return; } - String anchorId = context.getLast().getId(); + String anchorId = context.get(context.size() - 1).getId(); for (int i = msgs.size() - 1; i >= 0; i--) { if (anchorId.equals(msgs.get(i).getId())) { msgs.subList(0, i + 1).clear(); From 8ae9c596876d601cdfd04e478d440bff270226a7 Mon Sep 17 00:00:00 2001 From: xzxiaoshan <365384722@qq.com> Date: Mon, 3 Aug 2026 15:07:45 +0800 Subject: [PATCH 4/4] chore: trigger CI re-run for unrelated flaky HarnessAgentSubagentStreamEventsTest